Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3762aca307 | |||
| 4e27bde2d7 | |||
| d345b14e43 | |||
| 3ed7bfca67 | |||
| 6eb34a654f | |||
| 7084d99b89 | |||
| ae7845c375 | |||
| 61115f6f61 | |||
| 4b9ebda1b3 | |||
| 1e68d7ee39 | |||
| 6cccd458d4 | |||
| 42820fbe75 | |||
| d91ff886da | |||
| 386e760a5c |
@@ -24,6 +24,7 @@ import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
import dev.ltms.fleet.inject.TurnRegistrar;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
@@ -80,7 +81,9 @@ import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.BooleanSupplier;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Predicate;
|
||||
@@ -218,23 +221,23 @@ public final class Fleetd {
|
||||
// there is no 2-arg overload left for any lambda to silently bind to instead), but so a
|
||||
// test can call the exact same object this line builds, instead of asserting a copy of its
|
||||
// shape (round 3's lesson).
|
||||
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
|
||||
// fleetd #589 Group 1: extracted to forwardingExhaustionSink(...) below (see that method's
|
||||
// javadoc) so a dedicated test can prove this factory keeps reading the reference live,
|
||||
// rather than a rebuilt copy of its shape.
|
||||
ExhaustionSink forwardingExhaustionSink = forwardingExhaustionSink(exhaustionSinkRef);
|
||||
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
|
||||
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
|
||||
// unless opencode is the only kind configured.
|
||||
// fleetd #589 Group 2: extracted to claudeCodeLauncher(...)/openCodeLauncher(...) below (see
|
||||
// those methods' javadoc) so a dedicated test can prove the CB-596 memberCredentials policy
|
||||
// supplier is actually wired to each adapter, not silently replaced with `() -> null`.
|
||||
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new ClaudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), null, config::get));
|
||||
adapters.add(claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg, config));
|
||||
}
|
||||
if (!opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), config::get, forwardingExhaustionSink));
|
||||
adapters.add(openCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||
opencodeProfiles, cfg, config, forwardingExhaustionSink));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
@@ -408,12 +411,12 @@ public final class Fleetd {
|
||||
// LiveExhaustedPatterns's class doc for why this replaces the old compiled-once-at-startup
|
||||
// map. A profile with no exhaustedPattern simply returns null here, so its workers keep
|
||||
// today's completion-fallback behaviour unchanged.
|
||||
LiveExhaustedPatterns liveExhaustedPatterns = new LiveExhaustedPatterns(() -> config.get().profiles());
|
||||
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
|
||||
.orElse(null);
|
||||
// fleetd #589 Group 1: both extracted to liveExhaustedPatterns(...)/
|
||||
// exhaustedPatternLookup(...) below (see those methods' javadoc) — this is the worst
|
||||
// consequence in the whole #589 sweep: silently losing either wiring means a genuine
|
||||
// usage-limit refusal is handed back as a real completion instead of BACKEND_EXHAUSTED.
|
||||
LiveExhaustedPatterns liveExhaustedPatterns = liveExhaustedPatterns(config);
|
||||
ExhaustedPatternLookup exhaustedPatterns = exhaustedPatternLookup(sessions::roster, liveExhaustedPatterns);
|
||||
// The startup coverage line still reports the boot-time snapshot only — it is printed once,
|
||||
// here, and a reload no longer needs to change what it said; exhaustionDetectionArmed (via
|
||||
// liveExhaustedPatterns.armed, wired into quarantineSource below) is what stays live.
|
||||
@@ -461,11 +464,13 @@ public final class Fleetd {
|
||||
// method's javadoc for the full fleetd #175/#234/#446 history this used to carry inline —
|
||||
// so a dedicated test can drive the exact ExhaustionSink main() builds, not a hand-rebuilt
|
||||
// copy of its shape.
|
||||
ExhaustionSink exhaustionSink = exhaustionSink(sessions, config, quarantine,
|
||||
quarantineReasonByCredential, cfg);
|
||||
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
|
||||
// now that `sessions` exists to resolve target -> session -> profile.
|
||||
exhaustionSinkRef.set(exhaustionSink);
|
||||
// fleetd #589 Group 1: both statements (build + set) folded into publishExhaustionSink(...)
|
||||
// below (see that method's javadoc), so a test can prove the reference is actually
|
||||
// repointed at the real sink, not silently left at ExhaustionSink.none().
|
||||
ExhaustionSink exhaustionSink = publishExhaustionSink(exhaustionSinkRef, sessions, config,
|
||||
quarantine, quarantineReasonByCredential, cfg);
|
||||
// fleetd #201 Unit 5: the production BackendErrorSink needs `pushLoop` (built further below,
|
||||
// after `sessions`) to tell a lead about an incident or an unmapped target — the same
|
||||
// construction-order cycle `exhaustionSinkRef` breaks above, broken the same way: a mutable
|
||||
@@ -501,21 +506,26 @@ public final class Fleetd {
|
||||
// fleetd #556: registration is wired directly to `completion`, not folded into the
|
||||
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
|
||||
// listener) throwing, regardless of call order. See TurnRegistrar's javadoc.
|
||||
// fleetd #589 Group 3 (:505): extracted to turnRegistrar(...) — see FleetdTurnRegistrarWiringTest.
|
||||
Injector injector = new Injector(router, turnListener, deliverable,
|
||||
presence::forget, completion::register);
|
||||
presence::forget, turnRegistrar(completion));
|
||||
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start();
|
||||
|
||||
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
|
||||
// unusable), fleetd stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
|
||||
// connection, so keep the reference to close it in the ordered shutdown hook.
|
||||
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open);
|
||||
// fleetd #589 Group 3 (:512): opener extracted to replyInboxOpener() — see
|
||||
// FleetdReplyInboxOpenerWiringTest.
|
||||
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), replyInboxOpener());
|
||||
// CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate
|
||||
// broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator:
|
||||
// block this is null and every lead path below is simply not wired, which is exactly the
|
||||
// behaviour before this ticket. It owns a broker connection, so keep the reference for the
|
||||
// ordered shutdown hook.
|
||||
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open);
|
||||
// fleetd #589 Group 3 (:518): opener extracted to leadMailboxOpener() — see
|
||||
// FleetdLeadMailboxOpenerWiringTest.
|
||||
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), leadMailboxOpener());
|
||||
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
|
||||
// The pin also feeds CallerResolver below: a primary running inside a herdr pane would
|
||||
// otherwise resolve as a worker and be refused every orchestration tool.
|
||||
@@ -585,10 +595,12 @@ public final class Fleetd {
|
||||
// fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also
|
||||
// gets a wall-clock source to detect and correct for that freeze. Every other decision
|
||||
// in FleetHealthMonitor stays on the monotonic clock, unchanged.
|
||||
// fleetd #589 Group 3 (:591): failTarget extracted to healthFailTarget(...) — see
|
||||
// FleetdHealthFailTargetWiringTest.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
|
||||
cfg.health().intervalOrDefault(),
|
||||
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
|
||||
cfg.health().workingSuspectAfterOrDefault(), healthFailTarget(messages));
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
@@ -608,27 +620,11 @@ public final class Fleetd {
|
||||
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
|
||||
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
|
||||
// reached /metrics — the delegation was unresolvable and nothing said so.
|
||||
sessions.onRelease(detail -> {
|
||||
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
||||
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
||||
String reason = "the worker session was released before it replied";
|
||||
if (detail.worktreePath() != null) {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
});
|
||||
// fleetd #589 Group 3 (:611-631): the whole cleanup lambda extracted to releaseCleanup(...)
|
||||
// — see FleetdReleaseCleanupWiringTest, and MessageService.abandon's javadoc for the
|
||||
// documented incident (a torn-down worker's rendezvous waiter left open) this lambda exists
|
||||
// to prevent.
|
||||
sessions.onRelease(releaseCleanup(messages, replyInbox, primaryRegistry));
|
||||
|
||||
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
@@ -676,6 +672,10 @@ public final class Fleetd {
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
|
||||
// fleetd #602 gauge-wiring: threads each lead's configured configDir into the
|
||||
// context gauge — see leadConfigDirSource's own doc for why this, not a hardcoded
|
||||
// null, is what fleet_list's context row now reads.
|
||||
leadConfigDirSource(() -> config.get().profiles(), leaders),
|
||||
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
|
||||
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
|
||||
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
|
||||
@@ -1009,6 +1009,120 @@ public final class Fleetd {
|
||||
cfg.profiles()::keySet, System::nanoTime);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: the forwarding {@link ExhaustionSink} handed to the adapters built
|
||||
* before {@code sessions} exists (see the {@code exhaustionSinkRef}/{@code
|
||||
* forwardingExhaustionSink} locals in {@code main}, just above {@link #capacitySource}'s call
|
||||
* site). Before this ticket, {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} was
|
||||
* built inline — nothing a test could call directly, so a mutation swapping the supplier for a
|
||||
* hardcoded {@code () -> ExhaustionSink.none()} compiled clean and left the suite green: the
|
||||
* forwarder would silently stop reading the reference at all, and {@link
|
||||
* #publishExhaustionSink} repointing that reference later would have no effect.
|
||||
*
|
||||
* <p>Extracted the same way {@link #capacitySource}/{@link #loopHealthSource} were, so {@code
|
||||
* FleetdExhaustionSinkForwardingWiringTest} can call this factory directly with a real {@link
|
||||
* AtomicReference}, mutate the reference AFTER the forwarder is built, and prove the forwarder
|
||||
* still reads it live rather than a fixed target captured at construction time.
|
||||
*/
|
||||
static ExhaustionSink forwardingExhaustionSink(AtomicReference<ExhaustionSink> exhaustionSinkRef) {
|
||||
return ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: publish the real {@link ExhaustionSink} — built the same way {@link
|
||||
* #exhaustionSink} always was — into the forwarding reference {@link #forwardingExhaustionSink}
|
||||
* built above, replacing {@code main}'s previously untested two-statement sequence ({@code
|
||||
* ExhaustionSink exhaustionSink = exhaustionSink(...); exhaustionSinkRef.set(exhaustionSink);}).
|
||||
* Before this ticket, nothing proved the {@code .set(...)} call actually received the real sink
|
||||
* rather than a hardcoded {@code ExhaustionSink.none()} — the whole point of {@code
|
||||
* exhaustionSinkRef} existing (fleetd #175) is that {@link
|
||||
* dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check, built before {@code sessions}
|
||||
* exists, keeps working once this line runs; silently keeping the reference at {@code none()}
|
||||
* would mean that check permanently does nothing, with the full suite still green because no
|
||||
* existing test drives this exact call site.
|
||||
*
|
||||
* <p>Returns the built sink so {@code main} can still pass it to {@link CompletionResolver}'s
|
||||
* constructor at the same call site it already does, without building it twice.
|
||||
*/
|
||||
static ExhaustionSink publishExhaustionSink(AtomicReference<ExhaustionSink> exhaustionSinkRef,
|
||||
SessionManager sessions, ConfigRef config, BackendQuarantine quarantine,
|
||||
Map<String, String> quarantineReasonByCredential, FleetConfig cfg) {
|
||||
ExhaustionSink sink = exhaustionSink(sessions, config, quarantine, quarantineReasonByCredential, cfg);
|
||||
exhaustionSinkRef.set(sink);
|
||||
return sink;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: the live {@code exhaustedPattern} source (fleetd #446) {@link
|
||||
* CompletionResolver} enforces on, extracted out of {@code main} for the same reason {@link
|
||||
* #capacitySource} was. Before this ticket {@code new LiveExhaustedPatterns(() ->
|
||||
* config.get().profiles())} was built inline; replacing the supplier with a hardcoded {@code ()
|
||||
* -> Map.of()} compiled clean and left the suite green, meaning every profile's {@code
|
||||
* exhaustedPattern} would silently stop being recognised and a genuine usage-limit refusal
|
||||
* would be handed back as a real completion instead of {@code BACKEND_EXHAUSTED}.
|
||||
*/
|
||||
static LiveExhaustedPatterns liveExhaustedPatterns(ConfigRef config) {
|
||||
return new LiveExhaustedPatterns(() -> config.get().profiles());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: the {@link ExhaustedPatternLookup} {@link CompletionResolver} enforces
|
||||
* on, resolving a herdr {@code target} to its session's profile and then to that profile's live
|
||||
* {@link LiveExhaustedPatterns#patternFor}. Extracted out of {@code main} the same way {@link
|
||||
* #worktreeBranchLookup} was — same {@code Supplier<List<MemberSession>>} roster shape, same
|
||||
* reason: before this ticket the lambda was built inline, and replacing it with {@code target ->
|
||||
* null} (the exact shape of {@link ExhaustedPatternLookup#none()}) compiled clean and left the
|
||||
* suite green. This is the worst consequence in the whole #589 sweep (see the ticket): a
|
||||
* genuine usage-limit refusal would stop being classified as {@code BACKEND_EXHAUSTED} and
|
||||
* would be handed back to a waiting {@code fleet_send} as if it were real completed work.
|
||||
*
|
||||
* @param roster the live member roster, normally {@code sessions::roster}
|
||||
*/
|
||||
static ExhaustedPatternLookup exhaustedPatternLookup(Supplier<List<MemberSession>> roster,
|
||||
LiveExhaustedPatterns liveExhaustedPatterns) {
|
||||
return target -> roster.get().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
|
||||
.orElse(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: the production {@link ClaudeCodeLauncher} adapter, extracted out of
|
||||
* {@code main} the same way {@link #capacitySource} was. Before this ticket the constructor
|
||||
* call (11 arguments, including the CB-596 {@code memberCredentials} policy supplier) was built
|
||||
* inline; replacing the {@code () -> config.get().memberCredentials()} argument with {@code ()
|
||||
* -> null} compiled clean and left the suite green — {@code memberCredentials} is not {@code
|
||||
* null} itself (a lambda is never {@code null}), so {@link
|
||||
* dev.ltms.fleet.member.HerdrPeerLauncher#applyMemberCredentialPolicy} sees {@code
|
||||
* memberCredentials.get() == null} and silently shadows nothing, reopening the exact CB-592
|
||||
* exposure gap CB-596's policy closed. {@code FleetdClaudeCodeLauncherCredentialWiringTest}
|
||||
* calls this factory with a real {@link ConfigRef} carrying a {@code memberCredentials:} block
|
||||
* and proves a known-but-not-allowed name is actually shadowed on {@code spawn()}.
|
||||
*/
|
||||
static ClaudeCodeLauncher claudeCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
SubscriptionGuard guard, Map<String, FleetConfig.Profile> claudeProfiles, FleetConfig cfg,
|
||||
ConfigRef config) {
|
||||
return new ClaudeCodeLauncher(agents, spaces, guard, claudeProfiles, cfg.effectiveDefaultProfile(),
|
||||
System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(), () -> config.get().memberCredentials(), null, config::get);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: the production {@link OpenCodeLauncher} adapter, the {@code opencode}
|
||||
* counterpart to {@link #claudeCodeLauncher} above and extracted for the identical reason: the
|
||||
* same {@code () -> config.get().memberCredentials()} argument, reopening the same CB-592
|
||||
* exposure gap if silently replaced with {@code () -> null}.
|
||||
*/
|
||||
static OpenCodeLauncher openCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
Map<String, FleetConfig.Profile> opencodeProfiles, FleetConfig cfg, ConfigRef config,
|
||||
ExhaustionSink forwardingExhaustionSink) {
|
||||
return new OpenCodeLauncher(agents, spaces, opencodeProfiles, cfg.effectiveDefaultProfile(),
|
||||
System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(), () -> config.get().memberCredentials(), config::get,
|
||||
forwardingExhaustionSink);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #426: package-private factory for {@code fleet_list}'s {@code healthCoverage} source,
|
||||
* extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
@@ -1066,6 +1180,104 @@ public final class Fleetd {
|
||||
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :505}): package-private factory for the {@link Injector}'s
|
||||
* {@link TurnRegistrar}, extracted out of {@code main} for the same reason {@link
|
||||
* #loopHealthSource} was — before this ticket {@code completion::register} was an inline
|
||||
* argument to {@code new Injector(...)}, so nothing could pin it directly. Replacing it with
|
||||
* {@link TurnRegistrar#NOOP} compiles clean and leaves every existing test green: {@code
|
||||
* onDelivered}'s own {@code captureBaseline} does the identical {@code inFlight} check-and-put a
|
||||
* moment later on the ordinary path, so the two are indistinguishable once a turn's delivery
|
||||
* finishes normally. The gap {@link TurnRegistrar}'s own javadoc (fleetd #556) exists to close is
|
||||
* a {@code turnListener} callback throwing between the two — {@link
|
||||
* FleetdTurnRegistrarWiringTest} pins that {@code register} itself (not just {@code
|
||||
* captureBaseline}) makes a delivered turn's waiter resolvable.
|
||||
*/
|
||||
static TurnRegistrar turnRegistrar(CompletionResolver completion) {
|
||||
return completion::register;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :591}): package-private factory for {@link
|
||||
* FleetHealthMonitor}'s {@code failTarget} callback, extracted out of {@code main} for the same
|
||||
* reason {@link #loopHealthSource} was. Before this ticket {@code messages::abandon} was an
|
||||
* inline argument to {@code new FleetHealthMonitor(...)}; replacing it with a no-op {@code
|
||||
* BiConsumer} compiles clean and leaves every existing test green, and in production it means a
|
||||
* member found {@code GONE}/{@code NEVER_READY} never fails the ticket waiting on it — the
|
||||
* caller reports {@code PENDING} for the full 30-minute async timeout instead of the immediate,
|
||||
* accurate failure CB-580 exists to give it. {@link FleetdHealthFailTargetWiringTest} pins that
|
||||
* the returned callback actually reaches the real {@link MessageService#abandon}.
|
||||
*/
|
||||
static BiConsumer<String, String> healthFailTarget(MessageService messages) {
|
||||
return messages::abandon;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :611-631}): package-private factory for the whole {@link
|
||||
* SessionManager#onRelease} cleanup callback, extracted out of {@code main} for the same reason
|
||||
* {@link #loopHealthSource} was. Before this ticket this was an inline lambda built directly
|
||||
* inside {@code main}; replacing its body with a no-op {@code detail -> { }} compiles clean and
|
||||
* leaves every existing test green, and in production it is the exact incident {@link
|
||||
* MessageService#abandon}'s own javadoc documents: a torn-down worker's rendezvous waiter is
|
||||
* left open, so a blocking {@code fleet_send} keeps blocking and an async one reports {@code
|
||||
* PENDING} for a hardcoded thirty minutes on every {@code fleet_stop} and every idle-reap.
|
||||
*
|
||||
* <p>{@link FleetdReleaseCleanupWiringTest} pins all three collaborator calls this lambda makes
|
||||
* — {@code messages.abandon}, {@code replyInbox.release}, and {@code
|
||||
* primaryRegistry.forgetDelegation} — each already tested on its own ({@code MessageServiceTest},
|
||||
* {@code PrimaryRegistryTest}), but never before proven to actually be reached from here.
|
||||
*/
|
||||
static Consumer<SessionManager.ReleaseDetail> releaseCleanup(MessageService messages, ReplyInbox replyInbox,
|
||||
PrimaryRegistry primaryRegistry) {
|
||||
return detail -> {
|
||||
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
||||
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
||||
String reason = "the worker session was released before it replied";
|
||||
if (detail.worktreePath() != null) {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :512}): package-private factory for {@code selectReplyInbox}'s
|
||||
* production {@link AmqpOpener}, extracted out of {@code main} for the same reason {@link
|
||||
* #loopHealthSource} was. Before this ticket {@code AmqpReplyInbox::open} was an inline argument
|
||||
* to the {@code selectReplyInbox(...)} call; replacing it with {@code (uri, prefetch) -> new
|
||||
* InMemoryReplyInbox()} compiles clean and leaves every existing test green — {@link
|
||||
* FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own injected opener and
|
||||
* never sees what {@code main} actually passes. {@link FleetdReplyInboxOpenerWiringTest} pins
|
||||
* that this returns the real opener by pointing it at a guaranteed-closed local port and
|
||||
* asserting the real network attempt throws — the inert stub never attempts a connection at all.
|
||||
*/
|
||||
static AmqpOpener replyInboxOpener() {
|
||||
return AmqpReplyInbox::open;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :518}): package-private factory for {@code
|
||||
* openLeadMailbox}'s production {@link LeadMailboxOpener}, extracted out of {@code main} for the
|
||||
* same reason {@link #replyInboxOpener} was — same gap, same fix, the lead-coordination mailbox
|
||||
* instead of the reply inbox. {@link FleetdLeadMailboxOpenerWiringTest} pins that this returns
|
||||
* the real opener the same way.
|
||||
*/
|
||||
static LeadMailboxOpener leadMailboxOpener() {
|
||||
return LeadMailbox::open;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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).
|
||||
@@ -1368,6 +1580,58 @@ public final class Fleetd {
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #602 gauge-wiring: per-lead-name factory for {@link FleetMcp.LeadConfigDirSource} — the
|
||||
* {@code CLAUDE_CONFIG_DIR} the named lead's own profile runs on, so {@code LeadContextGauge}
|
||||
* reads the transcript directory that lead's Claude Code process actually writes to, not always
|
||||
* the built-in {@code <user.home>/.claude} default.
|
||||
*
|
||||
* <p>Follows the same {@code fleet.leaders.<name>.profile} link {@link #leadSeatLookup} already
|
||||
* uses to find a lead's profile, one step further to that profile's own {@code configDir:}. A
|
||||
* lead entry that names no {@code profile:}, or whose named profile is not configured, or whose
|
||||
* profile sets no {@code configDir:} override, returns {@code null} — {@code LeadContextGauge}
|
||||
* then falls back to its own default, exactly as before this ticket.
|
||||
*
|
||||
* @param profiles the live profile map, normally {@code () -> config.get().profiles()} in
|
||||
* {@code main} — read live, like every other {@code configDir} lookup, so a
|
||||
* config reload takes effect on the very next {@code fleet_list} call
|
||||
* @param leaders {@code fleet.leaders}, read once at startup like {@link #leadSeatLookup}'s own
|
||||
* {@code leaders} parameter — passed as a plain map, never re-read from
|
||||
* {@code config.get()}
|
||||
*/
|
||||
static Function<String, String> leadConfigDirLookup(Supplier<Map<String, FleetConfig.Profile>> profiles,
|
||||
Map<String, FleetConfig.Leader> leaders) {
|
||||
return leadName -> {
|
||||
FleetConfig.Leader lead = leaders.get(leadName);
|
||||
if (lead == null || lead.profile() == null || lead.profile().isBlank()) {
|
||||
return null;
|
||||
}
|
||||
FleetConfig.Profile leadProfile = profiles.get().get(lead.profile());
|
||||
return leadProfile == null ? null : leadProfile.configDir();
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #602 gauge-wiring follow-up (PR #606 review comment 17353): {@code main} used to build
|
||||
* {@code new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(...))} inline, with nothing a test
|
||||
* could call directly. Measured on that shape: replacing the whole expression with {@code
|
||||
* FleetMcp.LeadConfigDirSource.none()} at the call site compiled with 0 errors and left the full
|
||||
* 1822-test suite green — the daemon could be changed to always report every lead's context as
|
||||
* {@code UNKNOWN}, forever, while every test stayed green. That is the same hand-built-vs-wired
|
||||
* shape as fleetd #561/#248/#426/#562 ({@link #loopHealthSource}).
|
||||
*
|
||||
* <p>The fix extracts the inline {@code new} into this factory, in the same style as {@link
|
||||
* #loopHealthSource}/{@link #capacitySource}/{@link #healthCoverageSource} — which is exactly
|
||||
* what makes it directly callable from {@code FleetdLeadConfigDirSourceWiringTest}. That test
|
||||
* calls this factory with real {@link FleetConfig.Profile}/{@link FleetConfig.Leader} fixtures and
|
||||
* asserts the returned source resolves a real {@code configDir} — a property that would be false
|
||||
* if this method's body were mutated to {@code return FleetMcp.LeadConfigDirSource.none();}.
|
||||
*/
|
||||
static FleetMcp.LeadConfigDirSource leadConfigDirSource(Supplier<Map<String, FleetConfig.Profile>> profiles,
|
||||
Map<String, FleetConfig.Leader> leaders) {
|
||||
return new FleetMcp.LeadConfigDirSource(leadConfigDirLookup(profiles, leaders));
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
|
||||
@@ -221,7 +221,8 @@ public final class CallerResolver {
|
||||
// The config/live binding names this pane as an architect slot's own. Same
|
||||
// unforgeable pane mapping; the live binding, never a request argument, decides.
|
||||
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
|
||||
// escalating a dev, hunter, or reviewer into an architect. Checked before the worker fallback.
|
||||
// escalating a dev, hunter or reviewer into an architect. Checked before
|
||||
// the worker fallback.
|
||||
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
|
||||
}
|
||||
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
|
||||
|
||||
@@ -42,7 +42,8 @@ public interface MemberLifecycle {
|
||||
* Try to bind a newly spawned {@code terminal} into the role it was granted.
|
||||
*
|
||||
* @return the role this session actually holds: {@code role} unchanged for a role with no
|
||||
* live slot-binding semantics (dev, hunter, reviewer), or when the bind succeeded; a fallback
|
||||
* live slot-binding semantics (dev, hunter, reviewer), or when the bind
|
||||
* succeeded; a fallback
|
||||
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
|
||||
* Callers must record THIS value on the session, never the requested {@code role}, so
|
||||
* a later roster read never reports a role the session does not hold (CB-619). In
|
||||
|
||||
@@ -20,8 +20,8 @@ import java.util.function.Supplier;
|
||||
*
|
||||
* <p>Two halves, split by who owns each:
|
||||
* <ul>
|
||||
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code hunters}/
|
||||
* {@code reviewers}
|
||||
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/
|
||||
* {@code hunters}/{@code reviewers}
|
||||
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
|
||||
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
|
||||
* re-reads {@code fleet:} on every call, through a supplier the same shape as
|
||||
|
||||
@@ -2059,8 +2059,9 @@ public record FleetConfig(
|
||||
"defaultProfile", "a role pool under 'fleet:' — an unqualified spawn now names a role,"
|
||||
+ " and that role's pool supplies the candidate profiles",
|
||||
"architects", "'fleet.architects'",
|
||||
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers' or"
|
||||
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
|
||||
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers',"
|
||||
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not"
|
||||
+ " a 'role:' field",
|
||||
"leaders", "'fleet.leaders'",
|
||||
"leadScan", "'fleet.leaders.<name>.tabPrefix' and '.scanIntervalSeconds' — lead"
|
||||
+ " discovery is now configured on the lead it discovers");
|
||||
|
||||
@@ -0,0 +1,307 @@
|
||||
package dev.ltms.fleet.lead;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.RandomAccessFile;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
/**
|
||||
* Reads how full a lead's own Claude Code context window is, from the transcript Claude Code
|
||||
* itself writes — never from the lead's pane (fleetd has {@code AgentControl.read} for that, and
|
||||
* this must not use it: a pane holds terminal text, not the structured usage numbers a transcript
|
||||
* carries, and scraping it would also race the lead's own rendering).
|
||||
*
|
||||
* <p><strong>Why this exists.</strong> A lead auto-compacts when its context fills — on the host
|
||||
* this was built for, that happened 30 times in one session, discarding roughly 250,000 tokens
|
||||
* and costing 46 seconds to 3 minutes each time, and fleetd had no way to see it coming. This
|
||||
* class is the first thing that looks.
|
||||
*
|
||||
* <p><strong>The route.</strong> Claude Code appends one JSON object per line to
|
||||
* {@code <configDir>/projects/<slug>/<sessionId>.jsonl}. {@code <slug>} is an undocumented,
|
||||
* internal encoding of the working directory — this class never derives it. Instead it lists the
|
||||
* one-level-deep subdirectories of {@code <configDir>/projects/} and looks for
|
||||
* {@code <sessionId>.jsonl} by name, so the slug rule can change without breaking this reader.
|
||||
*
|
||||
* <ul>
|
||||
* <li><strong>Live context</strong> is read off the last record in the read window that carries
|
||||
* a {@code message.usage} object: {@code input_tokens + cache_read_input_tokens +
|
||||
* cache_creation_input_tokens}. This is what actually fills the window — a plain
|
||||
* {@code input_tokens} count alone understates it once the conversation has any cached
|
||||
* prefix, which on a long-lived lead is always.</li>
|
||||
* <li><strong>Compaction history</strong> is a count of {@code subtype: "compact_boundary"}
|
||||
* records seen in the same read window — see {@link Reading#compactions()}. It is a count
|
||||
* within the window this reader actually looked at, not a lifetime total: a session with
|
||||
* more compactions than fit in {@link #TAIL_BYTES} of transcript will undercount. That
|
||||
* trade-off is deliberate — see {@link #TAIL_BYTES}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p><strong>Three states, not two (OK / HIGH / UNKNOWN).</strong> Every path that cannot
|
||||
* positively establish the live token count — a missing file, an unreadable one, a peer that
|
||||
* is not a Claude backend, or every line in the read window failing to parse as JSON — returns
|
||||
* {@link State#UNKNOWN} with no token number, never a default "0" or "ok" that would read as
|
||||
* "this lead is fine" when the honest answer is "I could not look".
|
||||
*
|
||||
* <p><strong>A torn final line does not mean UNKNOWN.</strong> {@code fleet_list} reads this
|
||||
* transcript while Claude Code may be mid-write on it, so the last line in the window can be cut
|
||||
* off mid-flush — that is an ordinary, expected race, not a sign the format has changed. Earlier
|
||||
* this class treated ANY unparseable last line as UNKNOWN, on the theory that "if the format
|
||||
* changes, we should see UNKNOWN". That reasoning does not hold: a real format change makes
|
||||
* <em>every</em> line in the window unparseable, not only the last one written. So a single
|
||||
* malformed line (most often the final, torn one, but the check is not position-specific) is
|
||||
* simply skipped rather than treated as fatal, and the reading is built from whatever lines in the
|
||||
* window did parse. Only when <em>none</em> of them parse — the real format-change signal — does
|
||||
* this return {@link State#UNKNOWN}, still with no stale number standing in for "I could not
|
||||
* tell".
|
||||
*
|
||||
* <p><strong>Bounded cost.</strong> {@code fleet_list} is polled constantly, so every read is
|
||||
* capped two ways: {@link #TAIL_BYTES} bounds how much of the transcript is ever read from disk
|
||||
* (never the whole 52 MB a long-lived transcript reaches on the host this was measured on), and
|
||||
* {@link #DEFAULT_CACHE_TTL_MILLIS} bounds how often that bounded read actually happens — a burst
|
||||
* of {@code fleet_list} calls inside one TTL window reads the file once. One instance's cache is
|
||||
* keyed by {@code (configDir, sessionId)}, so it is safe to share across every lead a single
|
||||
* {@code fleet_list} call reports on.
|
||||
*/
|
||||
public final class LeadContextGauge {
|
||||
|
||||
private static final String PROJECTS_DIR = "projects";
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
|
||||
/**
|
||||
* How many trailing bytes of a transcript a single read ever pulls off disk. Chosen so one
|
||||
* read comfortably spans many recent turns — each usage or compact_boundary record is at most
|
||||
* a few KB — while staying nowhere near the 52 MB a long session's real transcript reaches on
|
||||
* the host this was built for; reading that whole file on every {@code fleet_list} call is
|
||||
* exactly the cost this bound exists to avoid. 2 MiB holds on the order of hundreds of recent
|
||||
* lines even when a turn's tool output is unusually large, which is far more than needed to
|
||||
* find the most recent usage record and any recent compaction.
|
||||
*/
|
||||
static final int TAIL_BYTES = 2 * 1024 * 1024;
|
||||
|
||||
/**
|
||||
* How long a {@link Reading} is served from cache before the file is read again.
|
||||
* {@code fleet_list} is called constantly (by design — it is the fleet's own status probe), so
|
||||
* without a TTL a burst of calls would re-read the transcript tail once per call. 5 seconds is
|
||||
* short enough that a caller watching for a state change never waits long, and long enough that
|
||||
* a poll loop calling every second or two only touches disk once per window.
|
||||
*/
|
||||
static final long DEFAULT_CACHE_TTL_MILLIS = 5_000;
|
||||
|
||||
/**
|
||||
* Live tokens at or above this count report {@link State#HIGH}. On the host this was measured
|
||||
* on, auto-compaction actually fires around 267,000–270,000 tokens, but the point of a HIGH
|
||||
* state is to warn before that happens, not at it — 200,000 is the standard Claude context
|
||||
* window size and a sensible built-in default: no config key is required to pick it, and a
|
||||
* lead crossing it is already deep enough into its window that a compaction is foreseeable.
|
||||
*/
|
||||
static final long HIGH_THRESHOLD_TOKENS = 200_000;
|
||||
|
||||
/** The only peer kind this reader understands ({@code Agent.agentType()}'s wire value). */
|
||||
private static final String CLAUDE_AGENT_TYPE = "claude";
|
||||
|
||||
public enum State { OK, HIGH, UNKNOWN }
|
||||
|
||||
/**
|
||||
* @param state {@link State#UNKNOWN} whenever {@code tokens} could not be established
|
||||
* @param tokens live context tokens, or {@code null} exactly when {@code state} is
|
||||
* {@link State#UNKNOWN}
|
||||
* @param compactions {@code compact_boundary} records seen in the read window (see class
|
||||
* javadoc) — {@code 0} both for "genuinely none seen" and for "unknown",
|
||||
* since a caller that already sees {@code state: UNKNOWN} has no reason to
|
||||
* trust this number either way
|
||||
*/
|
||||
public record Reading(State state, Long tokens, int compactions) {
|
||||
static Reading unknown() {
|
||||
return new Reading(State.UNKNOWN, null, 0);
|
||||
}
|
||||
}
|
||||
|
||||
private record CacheEntry(Reading reading, long readAtMillis) {
|
||||
}
|
||||
|
||||
private final LongSupplier clock;
|
||||
private final long ttlMillis;
|
||||
private final Map<String, CacheEntry> cache = new ConcurrentHashMap<>();
|
||||
/** Test seam only (package-private) — counts real disk reads, i.e. cache misses. */
|
||||
private final AtomicInteger diskReads = new AtomicInteger();
|
||||
|
||||
public LeadContextGauge() {
|
||||
this(System::currentTimeMillis, DEFAULT_CACHE_TTL_MILLIS);
|
||||
}
|
||||
|
||||
/** Test seam: an injectable clock and TTL so cache expiry is provable without sleeping. */
|
||||
LeadContextGauge(LongSupplier clock, long ttlMillis) {
|
||||
this.clock = clock;
|
||||
this.ttlMillis = ttlMillis;
|
||||
}
|
||||
|
||||
/** How many times this instance has actually read a transcript off disk — test seam only. */
|
||||
int diskReadCount() {
|
||||
return diskReads.get();
|
||||
}
|
||||
|
||||
/**
|
||||
* @param configDir the lead's {@code CLAUDE_CONFIG_DIR}, or {@code null}/blank to use the
|
||||
* default {@code <user.home>/.claude} — the right answer for the common case
|
||||
* where the lead's profile sets no {@code configDir} override
|
||||
* @param sessionId the lead's own Claude session id ({@code Agent.sessionId()}), or
|
||||
* {@code null} when herdr has not resolved one yet
|
||||
* @param agentType the detected peer kind ({@code Agent.agentType()}); anything other than
|
||||
* {@code "claude"} (including {@code null}, meaning undetected) reports
|
||||
* {@link State#UNKNOWN} — this reader only understands Claude Code's own
|
||||
* transcript format
|
||||
*/
|
||||
public Reading read(String configDir, String sessionId, String agentType) {
|
||||
if (sessionId == null || sessionId.isBlank()) {
|
||||
return Reading.unknown();
|
||||
}
|
||||
if (!CLAUDE_AGENT_TYPE.equalsIgnoreCase(agentType)) {
|
||||
return Reading.unknown();
|
||||
}
|
||||
String base = (configDir == null || configDir.isBlank())
|
||||
? System.getProperty("user.home") + "/.claude"
|
||||
: configDir;
|
||||
String cacheKey = base + '\u0000' + sessionId;
|
||||
long now = clock.getAsLong();
|
||||
CacheEntry cached = cache.get(cacheKey);
|
||||
if (cached != null && now - cached.readAtMillis() < ttlMillis) {
|
||||
return cached.reading();
|
||||
}
|
||||
Reading fresh = readUncached(base, sessionId);
|
||||
cache.put(cacheKey, new CacheEntry(fresh, now));
|
||||
return fresh;
|
||||
}
|
||||
|
||||
private Reading readUncached(String base, String sessionId) {
|
||||
diskReads.incrementAndGet();
|
||||
Path file = findTranscript(base, sessionId);
|
||||
if (file == null) {
|
||||
return Reading.unknown();
|
||||
}
|
||||
TailRead tail;
|
||||
try {
|
||||
tail = tailBytes(file, TAIL_BYTES);
|
||||
} catch (IOException e) {
|
||||
return Reading.unknown();
|
||||
}
|
||||
return parse(tail);
|
||||
}
|
||||
|
||||
/**
|
||||
* Finds {@code <sessionId>.jsonl} under {@code <base>/projects/}, one level deep — never by
|
||||
* deriving the slug directory from a working directory (see class javadoc). Bounded to a
|
||||
* single {@code list()} of {@code projects/} itself: it never recurses further, so the cost is
|
||||
* the number of project directories, not the size of any transcript inside them.
|
||||
*/
|
||||
private Path findTranscript(String base, String sessionId) {
|
||||
Path projectsDir = Path.of(base, PROJECTS_DIR);
|
||||
if (!Files.isDirectory(projectsDir)) {
|
||||
return null;
|
||||
}
|
||||
String filename = sessionId + ".jsonl";
|
||||
Path direct = projectsDir.resolve(filename);
|
||||
if (Files.isRegularFile(direct)) {
|
||||
return direct;
|
||||
}
|
||||
try (var children = Files.list(projectsDir)) {
|
||||
return children.filter(Files::isDirectory)
|
||||
.map(dir -> dir.resolve(filename))
|
||||
.filter(Files::isRegularFile)
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
} catch (IOException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Package-private (not {@code private}): {@link #tailBytes} is a test seam, see its javadoc. */
|
||||
record TailRead(byte[] bytes, boolean fromStart) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Reads at most {@code maxBytes} trailing bytes of {@code file}. Package-private (not
|
||||
* {@code private}) so a test can assert directly on the returned array's length — "bytes
|
||||
* actually read", not on any parsed answer — without needing a file anywhere near
|
||||
* {@link #TAIL_BYTES} in size to prove the cap holds.
|
||||
*/
|
||||
static TailRead tailBytes(Path file, int maxBytes) throws IOException {
|
||||
try (RandomAccessFile raf = new RandomAccessFile(file.toFile(), "r")) {
|
||||
long length = raf.length();
|
||||
long start = Math.max(0, length - maxBytes);
|
||||
raf.seek(start);
|
||||
byte[] buf = new byte[(int) (length - start)];
|
||||
raf.readFully(buf);
|
||||
return new TailRead(buf, start == 0);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Parses the tail into a {@link Reading}. The first line is dropped unconditionally whenever
|
||||
* the tail is not the whole file (it starts mid-line, cut by {@link #TAIL_BYTES} — an expected
|
||||
* artefact of the bound, not a data problem). Every remaining line is then parsed on a
|
||||
* best-effort basis: a line that fails to parse (most often the last one, torn by a write this
|
||||
* read raced — see the "torn final line" section of the class javadoc) is skipped, not fatal.
|
||||
* Only when none of the remaining lines parse does this report {@link State#UNKNOWN}.
|
||||
*/
|
||||
private Reading parse(TailRead tail) {
|
||||
String text = new String(tail.bytes(), StandardCharsets.UTF_8);
|
||||
List<String> lines = new ArrayList<>(List.of(text.split("\n", -1)));
|
||||
if (!lines.isEmpty() && lines.get(lines.size() - 1).isEmpty()) {
|
||||
lines.remove(lines.size() - 1); // trailing newline leaves a phantom empty last element
|
||||
}
|
||||
if (!tail.fromStart() && !lines.isEmpty()) {
|
||||
lines.remove(0); // first line is a fragment cut by our own tail bound, not real data
|
||||
}
|
||||
if (lines.isEmpty()) {
|
||||
return Reading.unknown();
|
||||
}
|
||||
Long tokens = null;
|
||||
int compactions = 0;
|
||||
boolean anyLineParsed = false;
|
||||
for (String line : lines) {
|
||||
JsonNode node = tryParse(line);
|
||||
if (node == null) {
|
||||
// A malformed line — typically the last one, cut mid-flush by a write this read
|
||||
// raced — is skipped rather than treated as fatal. See the class javadoc's "torn
|
||||
// final line" section for why: a real format change makes EVERY line unparseable,
|
||||
// not only this one, and that case is still caught below by anyLineParsed.
|
||||
continue;
|
||||
}
|
||||
anyLineParsed = true;
|
||||
JsonNode usage = node.path("message").path("usage");
|
||||
if (usage.isObject()) {
|
||||
tokens = usage.path("input_tokens").asLong(0)
|
||||
+ usage.path("cache_read_input_tokens").asLong(0)
|
||||
+ usage.path("cache_creation_input_tokens").asLong(0);
|
||||
}
|
||||
if ("compact_boundary".equals(node.path("subtype").asText(null))) {
|
||||
compactions++;
|
||||
}
|
||||
}
|
||||
if (!anyLineParsed) {
|
||||
return Reading.unknown();
|
||||
}
|
||||
if (tokens == null) {
|
||||
return new Reading(State.UNKNOWN, null, compactions);
|
||||
}
|
||||
State state = tokens >= HIGH_THRESHOLD_TOKENS ? State.HIGH : State.OK;
|
||||
return new Reading(state, tokens, compactions);
|
||||
}
|
||||
|
||||
private JsonNode tryParse(String line) {
|
||||
try {
|
||||
return MAPPER.readTree(line);
|
||||
} catch (IOException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -12,6 +12,7 @@ import dev.ltms.fleet.metrics.Metrics;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.lead.LeadRollover;
|
||||
import dev.ltms.fleet.msg.LeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadMessage;
|
||||
@@ -115,6 +116,8 @@ public final class FleetMcp {
|
||||
private final OutageSource outage;
|
||||
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
|
||||
private final LeadSeatSource leadSeats;
|
||||
/** fleetd #602 gauge-wiring: see {@link LeadConfigDirSource}. */
|
||||
private final LeadConfigDirSource leadConfigDirs;
|
||||
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
|
||||
private final LeadChannel leadChannel;
|
||||
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
|
||||
@@ -126,6 +129,17 @@ public final class FleetMcp {
|
||||
* clean {@code NOT_CONFIGURED} refusal rather than throwing. See {@link #handover}.
|
||||
*/
|
||||
private final LeadRollover leadRollover;
|
||||
/**
|
||||
* "Lead context gauge": how full each lead's own Claude Code context window is, reported on
|
||||
* {@code fleet_list}'s {@code leads} rows (see {@link #contextView}). Built unconditionally, in
|
||||
* the field initializer rather than a constructor parameter — this is not a togglable feature
|
||||
* with an on/off config knob the way {@link OutageSource}/{@link LeadSeatSource} are: it needs
|
||||
* no config at all (see {@link LeadContextGauge}'s own javadoc for the built-in defaults), so
|
||||
* there is no "off" value to thread through every existing constructor call site. One instance
|
||||
* per daemon so its read cache (keyed by session, TTL'd) is actually shared across
|
||||
* {@code fleet_list} calls rather than rebuilt — and therefore useless — on every call.
|
||||
*/
|
||||
private final LeadContextGauge leadContextGauge = new LeadContextGauge();
|
||||
|
||||
/** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */
|
||||
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
|
||||
@@ -246,6 +260,29 @@ public final class FleetMcp {
|
||||
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #602 gauge-wiring: a lead name's configured {@code CLAUDE_CONFIG_DIR} override, fed to
|
||||
* {@link LeadContextGauge#read} so {@code fleet_list}'s {@code context} row reads the transcript
|
||||
* directory the lead's own profile actually writes to, not always the built-in
|
||||
* {@code <user.home>/.claude} default.
|
||||
*
|
||||
* <p>The same idiom as {@link LeadSeatSource} — a {@code FleetMcp} constructor field, not a
|
||||
* lookup {@code contextView} performs itself, because {@code FleetMcp} holds no
|
||||
* {@link dev.ltms.fleet.config.FleetConfig} and {@code contextView} is {@code static}. See
|
||||
* {@code Fleetd.leadConfigDirLookup} for the derivation: the same
|
||||
* {@code fleet.leaders.<name>.profile} link {@link LeadSeatSource} already follows, resolved to
|
||||
* that profile's own {@code configDir:}.
|
||||
*
|
||||
* @param configDirFor lead name → {@code configDir}, or {@code null} when the lead's entry names
|
||||
* no profile, or that profile sets no {@code configDir} override — either
|
||||
* way {@link LeadContextGauge} then falls back to its own built-in default,
|
||||
* exactly as before this ticket
|
||||
*/
|
||||
public record LeadConfigDirSource(Function<String, String> configDirFor) {
|
||||
/** Inert source — every lead reads {@link LeadContextGauge}'s built-in default {@code configDir}. */
|
||||
public static LeadConfigDirSource none() { return new LeadConfigDirSource(_ -> null); }
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
|
||||
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
|
||||
@@ -346,7 +383,7 @@ public final class FleetMcp {
|
||||
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
|
||||
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
|
||||
peers, leadRollover);
|
||||
LeadConfigDirSource.none(), peers, leadRollover);
|
||||
}
|
||||
|
||||
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
@@ -354,7 +391,8 @@ public final class FleetMcp {
|
||||
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
||||
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||
LeadSeatSource leadSeats, LeadConfigDirSource leadConfigDirs, List<String> peers,
|
||||
LeadRollover leadRollover) {
|
||||
Objects.requireNonNull(callers, "callers");
|
||||
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
||||
== AuthorizationMode.ENFORCED;
|
||||
@@ -364,6 +402,7 @@ public final class FleetMcp {
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
this.outage = Objects.requireNonNull(outage, "outage");
|
||||
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
|
||||
this.leadConfigDirs = Objects.requireNonNull(leadConfigDirs, "leadConfigDirs");
|
||||
this.healthCoverage = healthCoverage;
|
||||
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
|
||||
this.leadRollover = leadRollover;
|
||||
@@ -498,7 +537,7 @@ public final class FleetMcp {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
|
||||
leadSeats, callers.leads(),
|
||||
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
coordinatorVisibleTo(principal(exchange)));
|
||||
@@ -1601,7 +1640,7 @@ public final class FleetMcp {
|
||||
Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
|
||||
OutageSource.none(), LeadSeatSource.none(), leads, selfTerm, coordination, false);
|
||||
OutageSource.none(), LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1610,7 +1649,7 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none(), false);
|
||||
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, CoordinationSource.none(), false);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1627,7 +1666,7 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination, false);
|
||||
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1636,7 +1675,7 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination, false);
|
||||
LeadSeatSource.none(), new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1658,7 +1697,7 @@ public final class FleetMcp {
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
|
||||
leadSeats, leads, selfTerm, coordination, false);
|
||||
leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1682,14 +1721,24 @@ public final class FleetMcp {
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
|
||||
outage, leadSeats, leads, selfTerm, coordination, callerIsPrimary);
|
||||
outage, leadSeats, new LeadContextGauge(), LeadConfigDirSource.none(), leads, selfTerm, coordination, callerIsPrimary);
|
||||
}
|
||||
|
||||
/**
|
||||
* The canonical implementation. {@code contextGauge} is the "lead context gauge" (see
|
||||
* {@link LeadContextGauge}) — every wrapper overload above passes a freshly constructed one,
|
||||
* which is correct for them (none of them exercise repeated calls where a shared cache would
|
||||
* matter); the one caller that matters for caching, {@code fleet_list}'s MCP handler, passes
|
||||
* its own single long-lived instance instead (see {@code FleetMcp}'s {@code leadContextGauge}
|
||||
* field).
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
LoopHealthSource loopHealth,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
LeadSeatSource leadSeats, LeadContextGauge contextGauge,
|
||||
LeadConfigDirSource leadConfigDirs,
|
||||
Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
@@ -1698,7 +1747,8 @@ public final class FleetMcp {
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
List<Map<String, Object>> leadRows = leads.entrySet().stream()
|
||||
.sorted(Map.Entry.comparingByValue())
|
||||
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
|
||||
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm, contextGauge,
|
||||
leadConfigDirs))
|
||||
.toList();
|
||||
// fleetd #209: this is the caller-driven fleet_list read that actually reports
|
||||
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
|
||||
@@ -1993,15 +2043,25 @@ public final class FleetMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* One lead's row: its address, its name, and whether it can be reached right now.
|
||||
* One lead's row: its address, its name, whether it can be reached right now, and how full its
|
||||
* own Claude Code context window is.
|
||||
*
|
||||
* <p>{@code status} is herdr's live view, and {@code unknown} when herdr is not tracking that
|
||||
* pane as an agent — the honest answer, and the one that matters: a lead whose pane herdr cannot
|
||||
* see is a lead a {@code fleet_send} cannot be typed into. It is reported rather than hidden,
|
||||
* because a peer that has gone unreachable is exactly what the sender needs to know.
|
||||
*
|
||||
* <p>{@code context} is the "lead context gauge" (fleetd's context-usage visibility ticket):
|
||||
* {@code {state: "ok"|"high"|"unknown", tokens?: number, compactions: number}}, read from the
|
||||
* lead's transcript — see {@link LeadContextGauge}. Kept small on purpose (a token count, a
|
||||
* state, and a compaction count) rather than echoing the whole reading history: this is a
|
||||
* roster row a caller glances at, not a diagnostics dump. {@code tokens} is present only when
|
||||
* {@code state} is not {@code "unknown"} — never a stale or default number standing in for "I
|
||||
* could not tell".
|
||||
*/
|
||||
private static Map<String, Object> leadView(String terminal, String name, Agent live,
|
||||
String selfTerm) {
|
||||
String selfTerm, LeadContextGauge contextGauge,
|
||||
LeadConfigDirSource leadConfigDirs) {
|
||||
Map<String, Object> m = new LinkedHashMap<>();
|
||||
m.put("sessionId", terminal);
|
||||
m.put("name", name);
|
||||
@@ -2010,9 +2070,32 @@ public final class FleetMcp {
|
||||
if (terminal.equals(selfTerm)) {
|
||||
m.put("self", true);
|
||||
}
|
||||
String configDir = leadConfigDirs.configDirFor().apply(name);
|
||||
m.put("context", contextView(contextGauge, live, configDir));
|
||||
return m;
|
||||
}
|
||||
|
||||
/**
|
||||
* Reads the lead context gauge for one lead. {@code configDir} is this lead's configured
|
||||
* {@code CLAUDE_CONFIG_DIR} override (see {@link LeadConfigDirSource}), derived from
|
||||
* {@code fleet.leaders.<name>.profile} → that profile's own {@code configDir:} — or {@code null}
|
||||
* when the lead's entry names no profile, or that profile sets no override, in which case
|
||||
* {@link LeadContextGauge#read} falls back to its own built-in default
|
||||
* ({@code <user.home>/.claude}).
|
||||
*/
|
||||
private static Map<String, Object> contextView(LeadContextGauge contextGauge, Agent live, String configDir) {
|
||||
String sessionId = live == null ? null : live.sessionId();
|
||||
String agentType = live == null ? null : live.agentType();
|
||||
LeadContextGauge.Reading reading = contextGauge.read(configDir, sessionId, agentType);
|
||||
Map<String, Object> c = new LinkedHashMap<>();
|
||||
c.put("state", reading.state().name().toLowerCase());
|
||||
if (reading.tokens() != null) {
|
||||
c.put("tokens", reading.tokens());
|
||||
}
|
||||
c.put("compactions", reading.compactions());
|
||||
return c;
|
||||
}
|
||||
|
||||
/** {@code fleet_stop}: tear a worker down by its pane id. */
|
||||
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
|
||||
if (isBlank(paneId)) {
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: {@link Fleetd#claudeCodeLauncher} is the factory that replaced {@code
|
||||
* main}'s inline {@code new ClaudeCodeLauncher(...)} call, whose 10th argument is the CB-596 {@code
|
||||
* memberCredentials} policy supplier ({@code () -> config.get().memberCredentials()}). Before this
|
||||
* ticket that argument was untestable wiring: replacing it with {@code () -> null} compiled with 0
|
||||
* errors and left every existing test green, since no existing test builds the exact object {@code
|
||||
* main} wires and then spawns it. {@code memberCredentials} being a lambda is never itself {@code
|
||||
* null}, so {@link dev.ltms.fleet.member.HerdrPeerLauncher#applyMemberCredentialPolicy} sees {@code
|
||||
* memberCredentials.get() == null} and silently shadows nothing — reopening the exact CB-592
|
||||
* exposure gap CB-596's policy closed (gitea issue #82).
|
||||
*
|
||||
* <p>This test drives the factory with a real {@link ConfigRef} carrying a {@code
|
||||
* memberCredentials:} block, spawns through the resulting launcher, and inspects what {@code
|
||||
* tab.create} actually carried — the same observable surface {@code ClaudeCodeLauncherTest}'s
|
||||
* {@code everyKnownNameNotAllowedIsShadowedWithTheSentinel} uses for the launcher's own credential
|
||||
* policy, applied here to prove {@code main}'s wiring reaches it.
|
||||
*/
|
||||
class FleetdClaudeCodeLauncherCredentialWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
ltms-local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: coder
|
||||
memberCredentials:
|
||||
policy: deny-by-default
|
||||
known:
|
||||
- GITEA_ACCESS_TOKEN
|
||||
""";
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, String> startEnv(FakeHerdr herdr) {
|
||||
return (Map<String, String>) ((Map<String, Object>) herdr.lastCall("tab.create").params()).get("env");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("main's memberCredentials wiring reaches ClaudeCodeLauncher: a known-but-not-allowed "
|
||||
+ "name is shadowed on spawn")
|
||||
void memberCredentialsWiringReachesClaudeCodeLauncher(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
ClaudeCodeLauncher launcher = Fleetd.claudeCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
cfg.profiles(), cfg, config);
|
||||
launcher.spawn();
|
||||
|
||||
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
|
||||
assertNotNull(shadowed,
|
||||
"GITEA_ACCESS_TOKEN is 'known' but not 'allow'-ed in the loaded config — it must be "
|
||||
+ "explicitly shadowed on spawn; replacing the memberCredentials supplier with "
|
||||
+ "() -> null at the Fleetd.claudeCodeLauncher call site must fail this "
|
||||
+ "assertion, since a null policy shadows nothing");
|
||||
assertFalse(shadowed.isBlank(), "the overlay value must be non-blank");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
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 org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#exhaustedPatternLookup} is the factory that replaced {@code
|
||||
* main}'s inline lambda — resolve a herdr {@code target} to its session's profile, then to that
|
||||
* profile's live {@link LiveExhaustedPatterns#patternFor}. Same shape as {@link
|
||||
* Fleetd#worktreeBranchLookup} (which {@code FleetdWorktreeBranchLookupTest} pins the same way).
|
||||
*
|
||||
* <p>Before this ticket the lambda was built inline in {@code main} and untestable: replacing it
|
||||
* with {@code target -> null} — the exact shape of {@link ExhaustedPatternLookup#none()} — compiled
|
||||
* with 0 errors and left every existing test green. Per the ticket, this is the worst consequence
|
||||
* in the whole #589 sweep: a genuine usage-limit refusal would stop being classified as {@code
|
||||
* BACKEND_EXHAUSTED} and would be handed back to a waiting {@code fleet_send} as if it were real
|
||||
* completed work.
|
||||
*/
|
||||
class FleetdExhaustedPatternLookupWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
exhaustedPattern: "usage limit"
|
||||
""";
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
private static LiveExhaustedPatterns liveExhaustedPatterns(Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
return new LiveExhaustedPatterns(() -> ConfigRef.fixed(cfg).get().profiles());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a known target resolves through its session's profile to that profile's live pattern")
|
||||
void knownTargetResolvesThroughItsProfile(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = liveExhaustedPatterns(dir);
|
||||
ExhaustedPatternLookup lookup = Fleetd.exhaustedPatternLookup(
|
||||
() -> List.of(session("term1", "terra")), patterns);
|
||||
|
||||
Pattern resolved = lookup.patternFor("term1");
|
||||
|
||||
assertNotNull(resolved,
|
||||
"the lookup must resolve term1 -> profile 'terra' -> LiveExhaustedPatterns.patternFor("
|
||||
+ "'terra') — replacing the lambda body with 'target -> null' at the "
|
||||
+ "Fleetd.exhaustedPatternLookup call site must fail this assertion");
|
||||
assertTrue(resolved.matcher("the usage limit has been reached").find());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unknown target resolves to null, not a thrown exception")
|
||||
void unknownTargetResolvesToNull(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = liveExhaustedPatterns(dir);
|
||||
ExhaustedPatternLookup lookup = Fleetd.exhaustedPatternLookup(
|
||||
() -> List.of(session("term1", "terra")), patterns);
|
||||
|
||||
assertNull(lookup.patternFor("term_stranger"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#forwardingExhaustionSink} is the factory that replaced
|
||||
* {@code main}'s inline {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} (fleetd #175's
|
||||
* construction-order break: the adapters need a sink before {@code sessions} exists to build the
|
||||
* real one). Before this ticket that call site was untestable wiring: replacing the supplier
|
||||
* argument with a hardcoded {@code () -> ExhaustionSink.none()} compiled with 0 errors and left
|
||||
* every existing test green, because no test builds the object {@code main} actually wires and
|
||||
* then mutates the reference afterward — every existing {@code ExhaustionSink.forwardingTo} caller
|
||||
* in this codebase reads and writes the SAME reference within one test, so a hardcoded-none supplier
|
||||
* and a correctly-forwarding one are indistinguishable to them.
|
||||
*
|
||||
* <p>This test builds the reference, builds the forwarder from it, and only THEN repoints the
|
||||
* reference at a spy sink — the discriminating order fleetd #175's whole design depends on
|
||||
* ({@code exhaustionSinkRef} starts at {@code none()} and is repointed once {@code sessions}
|
||||
* exists). A forwarder that captured a fixed target at construction time (the inert form) can never
|
||||
* see that later repoint.
|
||||
*/
|
||||
class FleetdExhaustionSinkForwardingWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("the forwarder reads the reference live: repointing it AFTER construction is honoured")
|
||||
void forwarderReadsTheReferenceLiveNotAFixedTargetCapturedAtConstruction() {
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
ExhaustionSink forwarder = Fleetd.forwardingExhaustionSink(exhaustionSinkRef);
|
||||
|
||||
AtomicBoolean spyCalled = new AtomicBoolean(false);
|
||||
exhaustionSinkRef.set((target, reason, profile) -> spyCalled.set(true));
|
||||
|
||||
forwarder.onExhausted("term_x", "usage limit reached", "terra");
|
||||
|
||||
assertTrue(spyCalled.get(),
|
||||
"forwardingExhaustionSink must delegate to whatever exhaustionSinkRef currently "
|
||||
+ "holds — hardcoding the supplier to () -> ExhaustionSink.none() at the "
|
||||
+ "Fleetd.forwardingExhaustionSink call site must fail this assertion, "
|
||||
+ "since the spy set into the reference after construction would never run");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("before any repoint, the forwarder is inert — it starts at none(), not a crash")
|
||||
void beforeAnyRepointTheForwarderIsInert() {
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
ExhaustionSink forwarder = Fleetd.forwardingExhaustionSink(exhaustionSinkRef);
|
||||
|
||||
AtomicBoolean spyCalled = new AtomicBoolean(false);
|
||||
forwarder.onExhausted("term_x", "usage limit reached", "terra");
|
||||
|
||||
assertFalse(spyCalled.get(), "nothing was ever wired to be called here — this only pins "
|
||||
+ "that the factory does not throw before a real sink is published");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#publishExhaustionSink} is the factory that replaced {@code
|
||||
* main}'s previously untested two-statement sequence — build the real {@link
|
||||
* Fleetd#exhaustionSink}, then {@code exhaustionSinkRef.set(exhaustionSink)}. {@link
|
||||
* Fleetd#exhaustionSink} itself is already pinned by {@code FleetdExhaustionSinkWarningTest} (its
|
||||
* log text) — what was NEVER pinned is the {@code .set(...)} call: {@code main} could replace it
|
||||
* with {@code exhaustionSinkRef.set(ExhaustionSink.none())} and compile with 0 errors, leaving
|
||||
* every existing test green, because {@link Fleetd#exhaustionSink}'s own tests build and call the
|
||||
* sink directly, never through the reference {@code main} publishes it into.
|
||||
*
|
||||
* <p>This test proves the PUBLISHED reference — not a freshly rebuilt sink — is the one that
|
||||
* actually quarantines a credential, by reading {@link BackendQuarantine#isQuarantined} after
|
||||
* calling {@code exhaustionSinkRef.get().onExhausted(...)}, the same object {@link
|
||||
* Fleetd#forwardingExhaustionSink} forwards to in production.
|
||||
*/
|
||||
class FleetdExhaustionSinkPublishWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""";
|
||||
|
||||
private static SessionManager emptyRosterSessions() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FleetConfig.Profile dummy = new FleetConfig.Profile(
|
||||
"dummy", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(dummy.profile(), dummy), dummy.profile(), _ -> "tok");
|
||||
// Never acquires a session — publishExhaustionSink's built sink resolves target -> profile
|
||||
// via the profileHint fallback (fleetd #234), exactly like OpenCodeLauncher's real call
|
||||
// site does, so this never needs a populated roster.
|
||||
return new SessionManager(launcher);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the published reference actually quarantines — not a rebuilt-but-never-set sink")
|
||||
void publishedReferenceActuallyQuarantines(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
Map<String, String> reasonByCredential = new HashMap<>();
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
|
||||
Fleetd.publishExhaustionSink(exhaustionSinkRef, emptyRosterSessions(), config, quarantine,
|
||||
reasonByCredential, cfg);
|
||||
exhaustionSinkRef.get().onExhausted("term_x", "The usage limit has been reached", "terra");
|
||||
|
||||
assertTrue(quarantine.isQuarantined("terra"),
|
||||
"publishExhaustionSink must repoint exhaustionSinkRef at the REAL sink — "
|
||||
+ "replacing the .set(...) call with exhaustionSinkRef.set(ExhaustionSink.none()) "
|
||||
+ "at the Fleetd.publishExhaustionSink call site must fail this assertion, "
|
||||
+ "since none()'s onExhausted does nothing");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("before publishing, the reference is still inert — no quarantine, no crash")
|
||||
void beforePublishingTheReferenceIsStillInert(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
|
||||
|
||||
exhaustionSinkRef.get().onExhausted("term_x", "The usage limit has been reached", "terra");
|
||||
|
||||
assertFalse(quarantine.isQuarantined("terra"),
|
||||
"nothing was published yet — this only pins the starting state the other test's "
|
||||
+ "assertion actually distinguishes from");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :591}. {@code Fleetd.main} wires {@link
|
||||
* dev.ltms.fleet.health.FleetHealthMonitor}'s {@code failTarget} callback with {@code
|
||||
* messages::abandon} — before this ticket that was an inline argument to {@code new
|
||||
* FleetHealthMonitor(...)}. Measured: replacing it with a no-op {@code BiConsumer} at the call
|
||||
* site compiles with 0 errors and leaves the full suite green, because nothing else in the tree
|
||||
* ever drives that specific constructor argument. In production it means a member the monitor
|
||||
* classifies {@code GONE}/{@code NEVER_READY} never has its pending ticket failed — the caller
|
||||
* keeps reporting {@code PENDING} for the full 30-minute async timeout instead of the immediate,
|
||||
* accurate failure CB-580 exists to give it.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#healthFailTarget} directly — never {@code FleetHealthMonitor}
|
||||
* or {@code main} — against a real {@link MessageService}, using the same {@code sendAsync} +
|
||||
* {@code poll} observable {@link MessageServiceTest} already relies on to pin {@code
|
||||
* MessageService.abandon} itself.
|
||||
*/
|
||||
class FleetdHealthFailTargetWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.healthFailTarget delegates to the real MessageService.abandon, not a no-op")
|
||||
void healthFailTargetDelegatesToMessagesAbandon() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
|
||||
BiConsumer<String, String> failTarget = Fleetd.healthFailTarget(messages);
|
||||
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting(rendezvous);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
failTarget.accept(T, "member unreachable (health monitor)");
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) {
|
||||
break;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
|
||||
"Fleetd.healthFailTarget(messages) must return messages::abandon — replacing it "
|
||||
+ "with a no-op BiConsumer at the Fleetd.healthFailTarget call site means "
|
||||
+ "this ticket is never failed and keeps polling as PENDING");
|
||||
assertTrue(view.detail() != null && view.detail().contains("member unreachable"),
|
||||
"the failure reason passed to failTarget.accept must reach MessageService.abandon "
|
||||
+ "and end up in the ticket's detail");
|
||||
}
|
||||
|
||||
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,100 @@
|
||||
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.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #602 gauge-wiring: {@link Fleetd#leadConfigDirLookup} is the factory {@code Fleetd.main}
|
||||
* wires into {@code FleetMcp.LeadConfigDirSource} so {@code fleet_list}'s {@code context} row reads
|
||||
* the transcript directory a lead's OWN profile actually writes to, instead of always falling back
|
||||
* to the built-in {@code <user.home>/.claude} default (see {@code LeadContextGauge}).
|
||||
*
|
||||
* <p>{@code FleetMcpLeadContextGaugeWiringTest} proves the directory this factory returns is what
|
||||
* actually gets read; this class proves the factory's own matching logic — the same
|
||||
* {@code fleet.leaders.<name>.profile} link {@code Fleetd.leadSeatLookup} already follows (see
|
||||
* {@code FleetdLeadSeatLookupTest}), one step further to that profile's own {@code configDir:}.
|
||||
*/
|
||||
class FleetdLeadConfigDirLookupTest {
|
||||
|
||||
private static FleetConfig.Profile profileWithConfigDir(String name, String configDir) {
|
||||
return new FleetConfig.Profile(name, null, "claude-sonnet-5", configDir, null, null,
|
||||
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
|
||||
null, 3, true, null, null, 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 lead on a profile that sets configDir resolves to that directory")
|
||||
void leadOnAProfileWithConfigDirResolvesToIt() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
|
||||
|
||||
assertEquals("/mnt/opus-claude", lookup.apply("primary"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("changing the config from directory A to directory B changes what the lookup reports")
|
||||
void configChangeFromDirectoryAToDirectoryBChangesTheAnswer() {
|
||||
java.util.concurrent.atomic.AtomicReference<Map<String, FleetConfig.Profile>> profilesRef =
|
||||
new java.util.concurrent.atomic.AtomicReference<>(
|
||||
Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-a")));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
Function<String, String> lookup = Fleetd.leadConfigDirLookup(profilesRef::get, leaders);
|
||||
|
||||
assertEquals("/mnt/dir-a", lookup.apply("primary"), "must read directory A before the config changes");
|
||||
|
||||
profilesRef.set(Map.of("opus", profileWithConfigDir("opus", "/mnt/dir-b")));
|
||||
assertEquals("/mnt/dir-b", lookup.apply("primary"), "must read directory B once the LIVE config changes — "
|
||||
+ "a lookup that snapshotted the profile map at construction would still answer directory A here");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead entry with no `profile:` (recognise-only) resolves to null, not a thrown exception")
|
||||
void recogniseOnlyLeadWithNoProfileResolvesToNull() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
|
||||
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, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
|
||||
|
||||
assertNull(lookup.apply("primary"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead naming a profile that is not configured resolves to null, not a thrown exception")
|
||||
void leadOnAnUnconfiguredProfileResolvesToNull() {
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("ghost-profile"));
|
||||
Function<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, leaders);
|
||||
|
||||
assertNull(lookup.apply("primary"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead on a profile that sets no configDir override resolves to null")
|
||||
void leadOnAProfileWithNoConfigDirResolvesToNull() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
Function<String, String> lookup = Fleetd.leadConfigDirLookup(() -> profiles, leaders);
|
||||
|
||||
assertNull(lookup.apply("primary"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unrecognised lead name resolves to null, not a thrown exception")
|
||||
void unrecognisedLeadNameResolvesToNull() {
|
||||
Function<String, String> lookup = Fleetd.leadConfigDirLookup(Map::of, Map.of());
|
||||
|
||||
assertNull(lookup.apply("ghost-lead"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #602 gauge-wiring follow-up (PR #606 review comment 17353): {@code Fleetd.main}'s {@code
|
||||
* LeadConfigDirSource} local used to be a bare {@code new FleetMcp.LeadConfigDirSource(
|
||||
* leadConfigDirLookup(...))} built inline, with nothing a test could call directly. Measured on
|
||||
* that shape: replacing the whole expression with {@code FleetMcp.LeadConfigDirSource.none()} at
|
||||
* the call site compiled with 0 errors and left the full 1822-test suite green — the daemon could
|
||||
* be changed to always report every lead's context as {@code UNKNOWN}, forever, and no test would
|
||||
* notice. That is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426/#562
|
||||
* ({@code FleetdLoopHealthSourceWiringTest}).
|
||||
*
|
||||
* <p>{@code FleetMcpLeadContextGaugeWiringTest} and {@code FleetdLeadConfigDirLookupTest} both
|
||||
* predate this class and are both still correct — but neither can catch the mutation above. One
|
||||
* builds its own {@code FleetMcp} and hands it its own {@code LeadConfigDirSource}; the other
|
||||
* builds its own lookup and calls {@link Fleetd#leadConfigDirLookup} directly. Neither one ever
|
||||
* calls the thing {@code Fleetd.main} actually calls.
|
||||
*
|
||||
* <p>The fix extracts the inline {@code new} into {@link Fleetd#leadConfigDirSource}, a
|
||||
* package-private factory in the same style as {@link Fleetd#loopHealthSource}/{@link
|
||||
* Fleetd#capacitySource}/{@link Fleetd#healthCoverageSource} — which is exactly what makes it
|
||||
* directly callable here. This test calls that factory with real {@link FleetConfig.Profile}/
|
||||
* {@link FleetConfig.Leader} fixtures (the same shapes {@code FleetdLeadConfigDirLookupTest}
|
||||
* already uses) and asserts the returned source resolves a real {@code configDir} — a property
|
||||
* that would be false if {@link Fleetd#leadConfigDirSource} were mutated to {@code return
|
||||
* FleetMcp.LeadConfigDirSource.none();}. Measured: mutating exactly that line makes
|
||||
* {@link #resolvesTheRealConfiguredConfigDir()} fail ({@code expected: </mnt/opus-claude> but was:
|
||||
* <null>}); restoring it makes the whole suite green again.
|
||||
*
|
||||
* <p><b>What this class does not and cannot cover.</b> {@code main}'s own line —
|
||||
* {@code leadConfigDirSource(() -> config.get().profiles(), leaders)} — could itself be swapped
|
||||
* for a bare {@code FleetMcp.LeadConfigDirSource.none()}, bypassing this factory entirely. Measured:
|
||||
* that exact mutation compiles with 0 errors and leaves every test in this file, and the full
|
||||
* 1825-test suite, green. {@link Fleetd#loopHealthSource}'s own wiring test has the identical gap
|
||||
* for its own one-line call in {@code main} — no test in this codebase calls {@code Fleetd.main}
|
||||
* far enough to observe which factory call it made. This class narrows the gap from "nothing tests
|
||||
* the wiring" (the pre-extraction state this ticket found) to "the factory's own logic is pinned,
|
||||
* and main's call to it is a one-line, visually-verifiable delegation" — the same standard already
|
||||
* accepted for {@code loopHealthSource}/{@code capacitySource}/{@code healthCoverageSource}.
|
||||
*/
|
||||
class FleetdLeadConfigDirSourceWiringTest {
|
||||
|
||||
private static FleetConfig.Profile profileWithConfigDir(String name, String configDir) {
|
||||
return new FleetConfig.Profile(name, null, "claude-sonnet-5", configDir, null, null,
|
||||
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
|
||||
null, 3, true, null, null, null);
|
||||
}
|
||||
|
||||
private static FleetConfig.Leader leadOnProfile(String profile) {
|
||||
return new FleetConfig.Leader(profile, "lead: primary", 1, "lead:", 10, "claude", "claude-sonnet-5");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the returned source resolves the lead's REAL configured configDir, not a hardcoded null")
|
||||
void resolvesTheRealConfiguredConfigDir() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", "/mnt/opus-claude"));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
|
||||
FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(() -> profiles, leaders);
|
||||
|
||||
assertEquals("/mnt/opus-claude", source.configDirFor().apply("primary"),
|
||||
"the configDirFor function must delegate to the real leadConfigDirLookup — mutating "
|
||||
+ "Fleetd.leadConfigDirSource's own body to `return FleetMcp.LeadConfigDirSource.none();` "
|
||||
+ "must fail this assertion (measured: it does — see this test's class javadoc for the "
|
||||
+ "companion measurement on main's one-line call to this factory, which this assertion "
|
||||
+ "does not and structurally cannot cover)");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead on a profile with no configDir override still resolves to null, not a crash")
|
||||
void leadWithNoConfigDirOverrideResolvesToNull() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("opus", profileWithConfigDir("opus", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
|
||||
FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(() -> profiles, leaders);
|
||||
|
||||
assertNull(source.configDirFor().apply("primary"),
|
||||
"no configDir: override configured ⇒ null, so LeadContextGauge falls back to its own default");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unrecognised lead name resolves to null, not a thrown exception")
|
||||
void unrecognisedLeadNameResolvesToNull() {
|
||||
FleetMcp.LeadConfigDirSource source = Fleetd.leadConfigDirSource(Map::of, Map.of());
|
||||
|
||||
assertNull(source.configDirFor().apply("ghost-lead"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :518}. {@code Fleetd.main} passes {@code LeadMailbox::open} as
|
||||
* the {@link Fleetd.LeadMailboxOpener} argument to {@code openLeadMailbox(...)} — before this
|
||||
* ticket that method reference was inline at the call site. Measured: replacing it with the inert
|
||||
* {@code (uri, selfCoordId, prefetch) -> null} compiles with 0 errors and leaves the full suite
|
||||
* green — {@code FleetdLeadMailboxSelectionTest} drives {@code openLeadMailbox} with its own
|
||||
* injected opener and never observes what {@code main} itself actually passes.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#leadMailboxOpener} directly and proves it is the real,
|
||||
* network-attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must
|
||||
* throw, exactly mirroring {@link FleetdReplyInboxOpenerWiringTest} for the reply-inbox opener.
|
||||
* The inert form never attempts a connection and returns {@code null} without throwing, so it
|
||||
* fails this assertion silently.
|
||||
*/
|
||||
class FleetdLeadMailboxOpenerWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.leadMailboxOpener is the real LeadMailbox::open, not a stub that never connects")
|
||||
void leadMailboxOpenerAttemptsARealConnection() throws Exception {
|
||||
int closedPort;
|
||||
try (ServerSocket socket = new ServerSocket(0)) {
|
||||
closedPort = socket.getLocalPort();
|
||||
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
|
||||
|
||||
Fleetd.LeadMailboxOpener opener = Fleetd.leadMailboxOpener();
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/coord", "coord-1", 50),
|
||||
"Fleetd.leadMailboxOpener() must be LeadMailbox::open — a real network attempt "
|
||||
+ "against a genuinely unreachable broker must throw. The inert form "
|
||||
+ "(uri, selfCoordId, prefetch) -> null never attempts a connection and "
|
||||
+ "returns null instead of throwing, so it would fail this assertion "
|
||||
+ "silently.");
|
||||
assertTrue(thrown.getMessage().contains("cannot connect to AMQP coordination broker"),
|
||||
"must be LeadMailbox.open's own real failure message, not a different exception "
|
||||
+ "shape standing in for it");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 1: {@link Fleetd#liveExhaustedPatterns} is the factory that replaced {@code
|
||||
* main}'s inline {@code new LiveExhaustedPatterns(() -> config.get().profiles())}. Before this
|
||||
* ticket, that supplier argument was untestable wiring: replacing it with a hardcoded {@code () ->
|
||||
* Map.of()} compiled with 0 errors and left every existing test green, because {@code
|
||||
* LiveExhaustedPatternsTest} builds its own instance directly with a hand-supplied map and never
|
||||
* goes through {@code main}'s call site.
|
||||
*
|
||||
* <p>Silently losing this wiring means every profile's {@code exhaustedPattern} stops being
|
||||
* recognised — {@link Fleetd#exhaustedPatternLookup} would never see a match, and a genuine
|
||||
* usage-limit refusal would be handed back to a waiting {@code fleet_send} as real completed work.
|
||||
*/
|
||||
class FleetdLiveExhaustedPatternsWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
exhaustedPattern: "usage limit"
|
||||
gx:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
""";
|
||||
|
||||
private static ConfigRef loadConfig(Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
return ConfigRef.fixed(cfg);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile with a configured exhaustedPattern is armed, with a compiled matcher")
|
||||
void configuredProfileIsArmed(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = Fleetd.liveExhaustedPatterns(loadConfig(dir));
|
||||
|
||||
assertTrue(patterns.armed("terra"),
|
||||
"the config's live profiles() supplier must reach LiveExhaustedPatterns — hardcoding "
|
||||
+ "the supplier to () -> Map.of() at the Fleetd.liveExhaustedPatterns call "
|
||||
+ "site must fail this assertion");
|
||||
assertTrue(patterns.patternFor("terra").matcher("the usage limit has been reached").find());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile with no configured exhaustedPattern is not armed, but is still resolvable")
|
||||
void unconfiguredProfileIsNotArmed(@TempDir Path dir) throws Exception {
|
||||
LiveExhaustedPatterns patterns = Fleetd.liveExhaustedPatterns(loadConfig(dir));
|
||||
|
||||
assertFalse(patterns.armed("gx"), "'gx' has no exhaustedPattern configured");
|
||||
assertNull(patterns.patternFor("gx"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.member.OpenCodeLauncher;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 2: {@link Fleetd#openCodeLauncher} is the factory that replaced {@code main}'s
|
||||
* inline {@code new OpenCodeLauncher(...)} call — the {@code opencode} counterpart to {@link
|
||||
* Fleetd#claudeCodeLauncher}, extracted for the identical reason. Its {@code memberCredentials}
|
||||
* argument is the same {@code () -> config.get().memberCredentials()} supplier; replacing it with
|
||||
* {@code () -> null} compiled with 0 errors and left every existing test green before this ticket,
|
||||
* reopening the same CB-592 exposure gap CB-596's policy closed.
|
||||
*
|
||||
* <p>Same observable surface as {@code OpenCodeLauncherTest}'s own {@code memberCredentials} tests:
|
||||
* spawn through the launcher {@code main} actually wires and inspect what {@code tab.create}
|
||||
* carried.
|
||||
*/
|
||||
class FleetdOpenCodeLauncherCredentialWiringTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
gemini:
|
||||
kind: opencode
|
||||
model: google/gemini-2.5-pro
|
||||
memberCredentials:
|
||||
policy: deny-by-default
|
||||
known:
|
||||
- GITEA_ACCESS_TOKEN
|
||||
""";
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, String> startEnv(FakeHerdr herdr) {
|
||||
return (Map<String, String>) ((Map<String, Object>) herdr.lastCall("tab.create").params()).get("env");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("main's memberCredentials wiring reaches OpenCodeLauncher: a known-but-not-allowed "
|
||||
+ "name is shadowed on spawn")
|
||||
void memberCredentialsWiringReachesOpenCodeLauncher(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
OpenCodeLauncher launcher = Fleetd.openCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), cfg.profiles(), cfg, config, ExhaustionSink.none());
|
||||
launcher.spawn();
|
||||
|
||||
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
|
||||
assertNotNull(shadowed,
|
||||
"GITEA_ACCESS_TOKEN is 'known' but not 'allow'-ed in the loaded config — it must be "
|
||||
+ "explicitly shadowed on spawn; replacing the memberCredentials supplier with "
|
||||
+ "() -> null at the Fleetd.openCodeLauncher call site must fail this "
|
||||
+ "assertion, since a null policy shadows nothing");
|
||||
assertFalse(shadowed.isBlank(), "the overlay value must be non-blank");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :611-631}. {@code Fleetd.main} wires {@code
|
||||
* sessions.onRelease(...)} with a lambda that calls three collaborators — {@code
|
||||
* messages.abandon}, {@code replyInbox.release}, and {@code primaryRegistry.forgetDelegation} —
|
||||
* before this ticket built inline inside {@code main}. Measured: replacing the whole lambda body
|
||||
* with {@code detail -> { }} compiles with 0 errors and leaves the full suite green, because each
|
||||
* collaborator is separately tested in isolation ({@code MessageServiceTest}, {@code
|
||||
* PrimaryRegistryTest}) but nothing before this ticket drove the lambda that calls all three from
|
||||
* {@code main}.
|
||||
*
|
||||
* <p>{@code MessageService.abandon}'s own javadoc documents the consequence: without this,
|
||||
* tearing a worker down leaves its rendezvous waiter open, so a blocking {@code fleet_send} keeps
|
||||
* blocking and an async one reports {@code PENDING} for a hardcoded thirty minutes on every {@code
|
||||
* fleet_stop} and every idle-reap.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#releaseCleanup} directly — never {@code SessionManager} or
|
||||
* {@code main} — against real {@link MessageService}, {@link InMemoryReplyInbox}, and {@link
|
||||
* PrimaryRegistry} instances, and asserts each collaborator's own observable effect: the pending
|
||||
* async ticket transitions to {@code FAILED} (abandon), the inbox no longer owns the target's
|
||||
* queue (release), and the recorded delegation is forgotten (forgetDelegation).
|
||||
*/
|
||||
class FleetdReleaseCleanupWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.releaseCleanup reaches messages.abandon, replyInbox.release, and primaryRegistry.forgetDelegation")
|
||||
void releaseCleanupReachesAllThreeCollaborators() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
InMemoryReplyInbox replyInbox = new InMemoryReplyInbox();
|
||||
PrimaryRegistry primaryRegistry = new PrimaryRegistry(null);
|
||||
|
||||
// Set up the "before" state each collaborator's own effect is measured against.
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting(rendezvous);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
|
||||
"sanity: the async ticket is pending before cleanup runs");
|
||||
|
||||
replyInbox.own(T);
|
||||
replyInbox.publish(T, "msg-1", "hello");
|
||||
assertEquals(1, replyInbox.peek(T).size(),
|
||||
"sanity: the inbox owns T and holds one message before cleanup runs");
|
||||
|
||||
primaryRegistry.recordDelegation(T, "lead-1");
|
||||
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(T).orElse(null),
|
||||
"sanity: the delegation is recorded before cleanup runs");
|
||||
|
||||
Consumer<SessionManager.ReleaseDetail> cleanup =
|
||||
Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry);
|
||||
cleanup.accept(new SessionManager.ReleaseDetail(T, null, null, null, null));
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) {
|
||||
break;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
|
||||
"releaseCleanup must call messages.abandon(...) — an inert detail -> { } lambda "
|
||||
+ "leaves this ticket PENDING forever");
|
||||
assertTrue(view.detail() != null && view.detail().contains("released"),
|
||||
"the abandon reason must say the worker session was released");
|
||||
|
||||
assertTrue(replyInbox.peek(T).isEmpty(),
|
||||
"releaseCleanup must call replyInbox.release(...) — an inert lambda leaves the "
|
||||
+ "inbox still owning T with its message");
|
||||
|
||||
assertTrue(primaryRegistry.nudgeTargetFor(T).isEmpty(),
|
||||
"releaseCleanup must call primaryRegistry.forgetDelegation(...) — an inert lambda "
|
||||
+ "leaves the stale delegation in place");
|
||||
}
|
||||
|
||||
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :512}. {@code Fleetd.main} passes {@code AmqpReplyInbox::open}
|
||||
* as the {@link Fleetd.AmqpOpener} argument to {@code selectReplyInbox(...)} — before this ticket
|
||||
* that method reference was inline at the call site. Measured: replacing it with the inert {@code
|
||||
* (uri, prefetch) -> new InMemoryReplyInbox()} compiles with 0 errors and leaves the full suite
|
||||
* green — {@code FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own
|
||||
* injected opener (including one test that passes the real {@code AmqpReplyInbox::open}
|
||||
* explicitly) and never observes what {@code main} itself actually passes.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#replyInboxOpener} directly and proves it is the real, network-
|
||||
* attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must throw —
|
||||
* the same shape {@code FleetdReplyInboxSelectionTest.aRealUnreachableBrokerFallsBackViaTheRealOpener}
|
||||
* already relies on for {@code AmqpReplyInbox.open} itself. The inert form never attempts a
|
||||
* connection and never throws, so it fails this assertion silently (by returning normally).
|
||||
*/
|
||||
class FleetdReplyInboxOpenerWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.replyInboxOpener is the real AmqpReplyInbox::open, not a stub that never connects")
|
||||
void replyInboxOpenerAttemptsARealConnection() throws Exception {
|
||||
int closedPort;
|
||||
try (ServerSocket socket = new ServerSocket(0)) {
|
||||
closedPort = socket.getLocalPort();
|
||||
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
|
||||
|
||||
Fleetd.AmqpOpener opener = Fleetd.replyInboxOpener();
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/vh", 50),
|
||||
"Fleetd.replyInboxOpener() must be AmqpReplyInbox::open — a real network attempt "
|
||||
+ "against a genuinely unreachable broker must throw. The inert form "
|
||||
+ "(uri, prefetch) -> new InMemoryReplyInbox() never attempts a connection "
|
||||
+ "and never throws, so it would return normally here and fail this "
|
||||
+ "assertion silently.");
|
||||
assertTrue(thrown.getMessage().contains("cannot connect to AMQP broker"),
|
||||
"must be AmqpReplyInbox.open's own real failure message, not a different exception "
|
||||
+ "shape standing in for it");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.TurnRegistrar;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :505}. {@code Fleetd.main} wires the {@link
|
||||
* dev.ltms.fleet.inject.Injector}'s {@link TurnRegistrar} with {@code completion::register} —
|
||||
* before this ticket that was an inline argument to {@code new Injector(...)}, so nothing could
|
||||
* pin it directly. Measured: replacing it with {@link TurnRegistrar#NOOP} at the call site
|
||||
* compiles with 0 errors and leaves the full suite green, because {@code onDelivered}'s own {@code
|
||||
* captureBaseline} performs the identical {@code inFlight} check-and-put a moment later on the
|
||||
* ordinary delivery path — the two are indistinguishable unless something reads the resolver
|
||||
* between {@code register} and {@code onDelivered}, or {@code onDelivered} never runs at all (the
|
||||
* gap fleetd #556 introduced {@link TurnRegistrar} to close).
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#turnRegistrar} directly — never {@code Injector} or {@code
|
||||
* main} — and drives {@link CompletionResolver} entirely through its public API: {@link
|
||||
* TurnRegistrar#register} followed by {@link CompletionResolver#resolveBeforePostAction}, which
|
||||
* looks up the same {@code inFlight} entry {@code onTurnComplete} would. With the real registrar,
|
||||
* that entry exists and the waiter opened by {@link Rendezvous#open} resolves; with {@link
|
||||
* TurnRegistrar#NOOP} nothing was ever registered, {@code resolveBeforePostAction} finds no
|
||||
* in-flight turn, and the waiter is left exactly as it started — never done.
|
||||
*/
|
||||
class FleetdTurnRegistrarWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.turnRegistrar delegates to the real CompletionResolver, not a no-op")
|
||||
void turnRegistrarDelegatesToCompletionRegister() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ BUILD GREEN: 391 files\n❯ ");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
// fleetd#164: an ever-advancing fake clock stands in for the real time a turn would take
|
||||
// between delivery and resolution, so the MIN_TURN_NANOS "too fast" floor never trips here —
|
||||
// see MessageServiceTest's resolverClock for the same technique.
|
||||
AtomicLong clock = new AtomicLong();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
|
||||
() -> clock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
|
||||
|
||||
TurnRegistrar registrar = Fleetd.turnRegistrar(completion);
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
|
||||
registrar.register(T, new TurnToken(T, waiter));
|
||||
|
||||
// Mirrors what Injector.onStatus's confirmed working->idle boundary would trigger via
|
||||
// CompletionResolver.onTurnComplete — resolveBeforePostAction is the public, synchronous
|
||||
// twin of that path and reads the exact same inFlight entry register() must have written.
|
||||
completion.resolveBeforePostAction(T);
|
||||
|
||||
assertTrue(waiter.isDone(),
|
||||
"Fleetd.turnRegistrar(completion) must return completion::register — replacing it "
|
||||
+ "with TurnRegistrar.NOOP at the Fleetd.turnRegistrar call site means this "
|
||||
+ "turn is never registered with CompletionResolver, so resolveBeforePostAction "
|
||||
+ "finds no in-flight turn and this waiter is never resolved");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,238 @@
|
||||
package dev.ltms.fleet.lead;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.RandomAccessFile;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Ticket "lead context gauge" — fleetd could not see how full a lead's own Claude Code context
|
||||
* window was, so a lead auto-compacting was always a surprise. {@link LeadContextGauge} reads that
|
||||
* off the lead's own transcript. Five acceptance properties from the ticket, one test class each
|
||||
* (plus the three separate UNKNOWN cases the ticket calls out by name):
|
||||
* <ol>
|
||||
* <li>{@link #tokenCountTracksTheLastUsageRecordAndChangesWithIt()}</li>
|
||||
* <li>{@link #compactionCountTracksCompactBoundaryRecordsAndChangesWithIt()}</li>
|
||||
* <li>{@link #missingFileIsUnknown()}, {@link #unreadableFileIsUnknown()}</li>
|
||||
* <li>{@link #tornFinalLineFallsBackToTheLastGoodReading()},
|
||||
* {@link #everyLineUnparseableIsUnknown()} (fleetd #602 gauge-wiring, finding 2)</li>
|
||||
* <li>{@link #readNeverExceedsTheTailBound()}</li>
|
||||
* <li>{@link #secondReadInsideTtlDoesNotTouchDiskAgain()}</li>
|
||||
* </ol>
|
||||
*
|
||||
* <p>Every fixture lives under {@code @TempDir} — never the operator's real config directory (see
|
||||
* the ticket's hard constraint on this).
|
||||
*
|
||||
* <p><strong>fleetd #602 gauge-wiring, finding 2.</strong> The property above used to be "a
|
||||
* truncated or invalid LAST line reports UNKNOWN" — on the theory that a bad last line signals a
|
||||
* format change. That reasoning did not hold: {@code fleet_list} reads this transcript while Claude
|
||||
* Code may be mid-write on it, so a torn LAST line is an ordinary race, not a format change, and a
|
||||
* real format change makes EVERY line unparseable, not only the last one written. So the property
|
||||
* is now split in two: {@link #tornFinalLineFallsBackToTheLastGoodReading} (the torn-line case must
|
||||
* NOT destroy a good earlier reading) and its control, {@link #everyLineUnparseableIsUnknown} (only
|
||||
* when NOTHING in the window parses does this report UNKNOWN).
|
||||
*/
|
||||
class LeadContextGaugeTest {
|
||||
|
||||
private static final String SESSION_ID = "11111111-1111-1111-1111-111111111111";
|
||||
|
||||
/** One line Claude Code would write for a turn with the given live-context total. */
|
||||
private static String usageLine(long inputTokens, long cacheRead, long cacheCreation) {
|
||||
return "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{"
|
||||
+ "\"input_tokens\":" + inputTokens + ","
|
||||
+ "\"cache_read_input_tokens\":" + cacheRead + ","
|
||||
+ "\"cache_creation_input_tokens\":" + cacheCreation + "}}}";
|
||||
}
|
||||
|
||||
/** One line Claude Code writes when an auto-compaction happens. */
|
||||
private static String compactionLine() {
|
||||
return "{\"type\":\"system\",\"subtype\":\"compact_boundary\","
|
||||
+ "\"compactMetadata\":{\"preTokens\":250000,\"postTokens\":5000,"
|
||||
+ "\"trigger\":\"auto\",\"durationMs\":54000}}";
|
||||
}
|
||||
|
||||
/** Lays out {@code <configDir>/projects/<anySlug>/<sessionId>.jsonl} and writes {@code lines}. */
|
||||
private static String writeTranscript(Path configDir, String sessionId, String... lines) throws IOException {
|
||||
Path projectDir = configDir.resolve("projects").resolve("some-project-slug");
|
||||
Files.createDirectories(projectDir);
|
||||
Path file = projectDir.resolve(sessionId + ".jsonl");
|
||||
StringBuilder sb = new StringBuilder();
|
||||
for (String line : lines) {
|
||||
sb.append(line).append('\n');
|
||||
}
|
||||
Files.writeString(file, sb.toString(), StandardCharsets.UTF_8);
|
||||
return configDir.toString();
|
||||
}
|
||||
|
||||
// --- property 1: live token count tracks the last usage record -----------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("reported tokens equal the LAST usage record's total, and change when it does")
|
||||
void tokenCountTracksTheLastUsageRecordAndChangesWithIt(@TempDir Path tmp) throws IOException {
|
||||
String configDir = writeTranscript(tmp, SESSION_ID,
|
||||
usageLine(1_000, 0, 0),
|
||||
usageLine(40_000, 5_000, 3_000)); // last record: 48,000
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
|
||||
LeadContextGauge.Reading first = gauge.read(configDir, SESSION_ID, "claude");
|
||||
assertEquals(48_000L, first.tokens(), "must total input+cache_read+cache_creation of the LAST usage record");
|
||||
assertEquals(LeadContextGauge.State.OK, first.state());
|
||||
|
||||
// Change N in the fixture (a fresh session id avoids the cache) -- the reported number
|
||||
// must change with it, not stay pinned to the first fixture's total.
|
||||
String otherSession = "22222222-2222-2222-2222-222222222222";
|
||||
writeTranscript(tmp, otherSession, usageLine(100_000, 50_000, 50_000)); // last record: 200,000
|
||||
LeadContextGauge.Reading second = gauge.read(configDir, otherSession, "claude");
|
||||
assertEquals(200_000L, second.tokens());
|
||||
assertTrue(second.tokens() != first.tokens(), "changing N in the fixture must change the reported number");
|
||||
}
|
||||
|
||||
// --- property 2: compaction count tracks compact_boundary records --------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("reported compaction count equals K compact_boundary records, and changes with K")
|
||||
void compactionCountTracksCompactBoundaryRecordsAndChangesWithIt(@TempDir Path tmp) throws IOException {
|
||||
String sessionTwoCompactions = "33333333-3333-3333-3333-333333333333";
|
||||
writeTranscript(tmp, sessionTwoCompactions,
|
||||
usageLine(1_000, 0, 0),
|
||||
compactionLine(),
|
||||
usageLine(2_000, 0, 0),
|
||||
compactionLine(),
|
||||
usageLine(3_000, 0, 0));
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
LeadContextGauge.Reading twoCompactions = gauge.read(tmp.toString(), sessionTwoCompactions, "claude");
|
||||
assertEquals(2, twoCompactions.compactions());
|
||||
|
||||
String sessionZeroCompactions = "44444444-4444-4444-4444-444444444444";
|
||||
writeTranscript(tmp, sessionZeroCompactions, usageLine(3_000, 0, 0));
|
||||
LeadContextGauge.Reading zeroCompactions = gauge.read(tmp.toString(), sessionZeroCompactions, "claude");
|
||||
assertEquals(0, zeroCompactions.compactions(), "changing K in the fixture must change the reported count");
|
||||
}
|
||||
|
||||
// --- property 3: two separate UNKNOWN cases -------------------------------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("a missing transcript file reports UNKNOWN with no token number")
|
||||
void missingFileIsUnknown(@TempDir Path tmp) {
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
LeadContextGauge.Reading reading = gauge.read(tmp.toString(), SESSION_ID, "claude");
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state());
|
||||
assertNull(reading.tokens());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unreadable transcript file reports UNKNOWN with no token number")
|
||||
void unreadableFileIsUnknown(@TempDir Path tmp) throws IOException {
|
||||
String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0));
|
||||
Path file = tmp.resolve("projects").resolve("some-project-slug").resolve(SESSION_ID + ".jsonl");
|
||||
assertTrue(file.toFile().setReadable(false), "test setup: must be able to revoke read permission");
|
||||
try {
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude");
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state());
|
||||
assertNull(reading.tokens());
|
||||
} finally {
|
||||
file.toFile().setReadable(true); // so @TempDir cleanup can delete it
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("fleetd #602 finding 2: a torn final line does not destroy a good earlier reading")
|
||||
void tornFinalLineFallsBackToTheLastGoodReading(@TempDir Path tmp) throws IOException {
|
||||
// Deliberately NOT via writeTranscript: that helper appends a trailing newline after every
|
||||
// line, including the last one, which would not model what a write caught mid-flush looks
|
||||
// like on disk. Claude Code appends and flushes one line at a time, so the torn line here has
|
||||
// no trailing newline at all -- exactly the shape a read racing an in-progress write sees.
|
||||
Path projectDir = tmp.resolve("projects").resolve("some-project-slug");
|
||||
Files.createDirectories(projectDir);
|
||||
Path file = projectDir.resolve(SESSION_ID + ".jsonl");
|
||||
String lastCompleteLine = usageLine(1_000, 2_000, 3_000); // last COMPLETE record: 6,000
|
||||
String tornLine = "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{\"input_tok";
|
||||
Files.writeString(file, lastCompleteLine + "\n" + tornLine, StandardCharsets.UTF_8);
|
||||
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
LeadContextGauge.Reading reading = gauge.read(tmp.toString(), SESSION_ID, "claude");
|
||||
assertEquals(LeadContextGauge.State.OK, reading.state(),
|
||||
"a torn final line must not turn a good earlier reading into UNKNOWN");
|
||||
assertEquals(6_000L, reading.tokens(),
|
||||
"must report the last COMPLETE line's total, not fail the whole read");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("fleetd #602 finding 2 control: when EVERY line is unparseable, the reading IS UNKNOWN")
|
||||
void everyLineUnparseableIsUnknown(@TempDir Path tmp) throws IOException {
|
||||
// The real signal a format change gives: not one bad line, but ALL of them. Without this
|
||||
// control, code that never reports UNKNOWN at all would still pass the torn-line property.
|
||||
String configDir = writeTranscript(tmp, SESSION_ID,
|
||||
"{this is not json at all",
|
||||
"neither is this{{{");
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
LeadContextGauge.Reading reading = gauge.read(configDir, SESSION_ID, "claude");
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, reading.state(),
|
||||
"every line unparseable is the real format-change signal and must still report UNKNOWN");
|
||||
assertNull(reading.tokens());
|
||||
}
|
||||
|
||||
// --- property 4: the tail bound is actually enforced ----------------------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("reading a transcript far larger than the tail bound reads no more than the bound")
|
||||
void readNeverExceedsTheTailBound(@TempDir Path tmp) throws IOException {
|
||||
int smallBound = 4_096; // exercise the mechanism without writing a multi-MB fixture
|
||||
Path file = tmp.resolve("big.jsonl");
|
||||
// One line far bigger than smallBound, repeated, so the whole file is many times the bound.
|
||||
String line = usageLine(42, 0, 0) + " ".repeat(500);
|
||||
StringBuilder sb = new StringBuilder();
|
||||
for (int i = 0; i < 50; i++) {
|
||||
sb.append(line).append('\n');
|
||||
}
|
||||
Files.writeString(file, sb.toString(), StandardCharsets.UTF_8);
|
||||
long fileLength = Files.size(file);
|
||||
assertTrue(fileLength > (long) smallBound * 5, "test setup: fixture must genuinely dwarf the bound");
|
||||
|
||||
var tail = LeadContextGauge.tailBytes(file, smallBound);
|
||||
assertEquals(smallBound, tail.bytes().length,
|
||||
"must read exactly the bound, not the whole " + fileLength + "-byte file — assert on bytes actually read");
|
||||
}
|
||||
|
||||
// --- property 5: the cache TTL means a burst of calls reads the file once -------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("the same (configDir, sessionId) read twice inside the TTL touches disk once")
|
||||
void secondReadInsideTtlDoesNotTouchDiskAgain(@TempDir Path tmp) throws IOException {
|
||||
String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0));
|
||||
AtomicLong now = new AtomicLong(0);
|
||||
LeadContextGauge gauge = new LeadContextGauge(now::get, 5_000);
|
||||
|
||||
gauge.read(configDir, SESSION_ID, "claude");
|
||||
gauge.read(configDir, SESSION_ID, "claude"); // still inside the TTL window
|
||||
assertEquals(1, gauge.diskReadCount(), "two reads inside the TTL must touch disk once");
|
||||
|
||||
now.set(6_000); // past the TTL
|
||||
gauge.read(configDir, SESSION_ID, "claude");
|
||||
assertEquals(2, gauge.diskReadCount(), "a read past the TTL must touch disk again");
|
||||
}
|
||||
|
||||
// --- non-Claude peer / unresolved session -----------------------------------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("a non-claude agent type or a null session id is UNKNOWN, never OK")
|
||||
void nonClaudePeerOrUnresolvedSessionIsUnknown(@TempDir Path tmp) throws IOException {
|
||||
String configDir = writeTranscript(tmp, SESSION_ID, usageLine(1_000, 0, 0));
|
||||
LeadContextGauge gauge = new LeadContextGauge();
|
||||
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, SESSION_ID, "opencode").state());
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, SESSION_ID, null).state());
|
||||
assertEquals(LeadContextGauge.State.UNKNOWN, gauge.read(configDir, null, "claude").state());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,126 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #602 gauge-wiring: {@code FleetMcp}'s {@code fleet_list} handler shipped calling
|
||||
* {@code LeadContextGauge.read(null, ...)} unconditionally, so every lead whose profile sets a
|
||||
* {@code configDir:} override read the wrong transcript directory forever, with no error anywhere —
|
||||
* {@link LeadContextGauge}'s own 8 tests all passed because none of them exercised THIS wiring; each
|
||||
* one hands the gauge its own directory directly.
|
||||
*
|
||||
* <p>This class proves the directory {@code fleet.leaders.<name>.profile}'s own {@code configDir:}
|
||||
* names is the one {@code fleet_list}'s {@code context} row actually reads — not merely that SOME
|
||||
* directory got passed. If the wiring in {@code FleetMcp.leadView}/{@code contextView} is ever
|
||||
* reverted to a hardcoded {@code null}, {@link #configuredDirectoryDecidesWhichTranscriptIsRead()}
|
||||
* must go red: both directories below hold a REAL, DIFFERENT token count for the SAME session id,
|
||||
* so a hardcoded {@code null} (always reading the built-in default, where neither transcript lives)
|
||||
* would report {@code UNKNOWN} both times instead of the two distinct numbers this test asserts.
|
||||
*/
|
||||
class FleetMcpLeadContextGaugeWiringTest {
|
||||
|
||||
private static final String LEAD_TERMINAL = "term_lead";
|
||||
private static final String LEAD_NAME = "opus";
|
||||
/** {@code FakeHerdr.withAgent} always projects {@code agent_session.value} as {@code sess-<terminalId>}. */
|
||||
private static final String LEAD_SESSION_ID = "sess-" + LEAD_TERMINAL;
|
||||
|
||||
private static ClaudeCodeLauncher workerService(FakeHerdr h) {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> "tok");
|
||||
}
|
||||
|
||||
private static String textOf(McpSchema.CallToolResult r) {
|
||||
return ((McpSchema.TextContent) r.content().getFirst()).text();
|
||||
}
|
||||
|
||||
/** One line Claude Code would write for a turn with the given live-context total. */
|
||||
private static String usageLine(long tokens) {
|
||||
return "{\"type\":\"assistant\",\"message\":{\"role\":\"assistant\",\"usage\":{"
|
||||
+ "\"input_tokens\":" + tokens + ",\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}";
|
||||
}
|
||||
|
||||
/** Lays out {@code <configDir>/projects/<anySlug>/<sessionId>.jsonl} carrying one usage record. */
|
||||
private static String writeTranscript(Path configDir, String sessionId, long tokens) throws IOException {
|
||||
Path projectDir = configDir.resolve("projects").resolve("some-project-slug");
|
||||
Files.createDirectories(projectDir);
|
||||
Files.writeString(projectDir.resolve(sessionId + ".jsonl"), usageLine(tokens) + "\n", StandardCharsets.UTF_8);
|
||||
return configDir.toString();
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the CONFIGURED directory decides which transcript is read, and follows a config change")
|
||||
void configuredDirectoryDecidesWhichTranscriptIsRead(@TempDir Path tmp) throws IOException {
|
||||
Path dirA = Files.createDirectory(tmp.resolve("dir-a"));
|
||||
Path dirB = Files.createDirectory(tmp.resolve("dir-b"));
|
||||
writeTranscript(dirA, LEAD_SESSION_ID, 11_000);
|
||||
writeTranscript(dirB, LEAD_SESSION_ID, 22_000);
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("lead-opus", LEAD_TERMINAL, "wL:p1", "wL:t1");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr));
|
||||
LeadContextGauge contextGauge = new LeadContextGauge();
|
||||
AtomicReference<String> configuredDir = new AtomicReference<>(dirA.toString());
|
||||
FleetMcp.LeadConfigDirSource source = new FleetMcp.LeadConfigDirSource(name ->
|
||||
LEAD_NAME.equals(name) ? configuredDir.get() : null);
|
||||
|
||||
String firstRead = textOf(listFleet(herdr, sessions, contextGauge, source));
|
||||
assertTrue(firstRead.contains("\"tokens\":11000"),
|
||||
"config names directory A ⇒ fleet_list must report A's token count: " + firstRead);
|
||||
|
||||
// The config changes to name directory B instead -- the very next read must follow it.
|
||||
configuredDir.set(dirB.toString());
|
||||
String secondRead = textOf(listFleet(herdr, sessions, contextGauge, source));
|
||||
assertTrue(secondRead.contains("\"tokens\":22000"),
|
||||
"config now names directory B ⇒ fleet_list must report B's token count, not A's stale one: "
|
||||
+ secondRead);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead whose config names no directory still degrades to the built-in default, never throws")
|
||||
void aLeadWhoseConfigNamesNoDirectoryStillFallsBackWithoutThrowing() {
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("lead-opus", LEAD_TERMINAL, "wL:p1", "wL:t1");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr));
|
||||
LeadContextGauge contextGauge = new LeadContextGauge();
|
||||
|
||||
McpSchema.CallToolResult res = listFleet(herdr, sessions, contextGauge, FleetMcp.LeadConfigDirSource.none());
|
||||
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a lead with no configured configDir must degrade, never throw");
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"state\":\"unknown\""), "no transcript under the built-in default ⇒ UNKNOWN: " + out);
|
||||
assertFalse(out.contains("\"tokens\""), "UNKNOWN must never carry a stale/default token number: " + out);
|
||||
}
|
||||
|
||||
private static McpSchema.CallToolResult listFleet(FakeHerdr herdr, SessionManager sessions,
|
||||
LeadContextGauge contextGauge, FleetMcp.LeadConfigDirSource leadConfigDirs) {
|
||||
return FleetMcp.listFleet(workerService(herdr), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.LoopHealthSource.none(), FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), contextGauge, leadConfigDirs,
|
||||
Map.of(LEAD_TERMINAL, LEAD_NAME), LEAD_TERMINAL, FleetMcp.CoordinationSource.none(), false);
|
||||
}
|
||||
}
|
||||
@@ -169,7 +169,55 @@ hash256() {
|
||||
# on PATH). "absent" must never be the answer for a file that exists — that conflation, on Linux,
|
||||
# was the whole defect this ticket fixes.
|
||||
jar_id() { local f="${1:-$JAR}"; [ -f "$f" ] && hash256 "$f" || echo "absent"; }
|
||||
running_pid() { pgrep -f "$PATTERN" || true; }
|
||||
|
||||
# fleetd #593 — `pgrep -f "$PATTERN"` matches ANY process whose full command line CONTAINS the
|
||||
# pattern text, and that is not the same thing as "is the daemon". A shell that merely embeds the
|
||||
# pattern as literal text — a human typing this exact investigation by hand, an ssh-shaped
|
||||
# `sh -c '...; ...'`, a pipeline, or any other non-exec'ing shell that never replaced itself with
|
||||
# the pattern-holding command — still shows up in that match, and it is the INSTRUMENT, not the
|
||||
# daemon. Measured live on this Mac: `sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 30' &`
|
||||
# leaves a real `sh` process alive (it forks for the `sleep`, it does not exec into it) whose own
|
||||
# `ps -o args` is `sh -c echo "target/fleetd.jar" >/dev/null; sleep 30` — `pgrep -f "$PATTERN"`
|
||||
# matches that line right alongside the real `java -jar target/fleetd.jar` process. `pgrep -c`
|
||||
# (an in-one-call count) does not exist on BSD/macOS at all, so this cannot be fixed by switching
|
||||
# pgrep flags — it has to filter what pgrep already found, after the fact, in a way that still
|
||||
# runs on BSD.
|
||||
#
|
||||
# fleetd #593 CORRECTION 1 — the first cut of this filter kept everything whose `comm` was NOT a
|
||||
# shell name (a denylist: sh/bash/zsh/dash/ksh). Two holes in that, both the same false-positive
|
||||
# shape the ticket exists to remove in the first place:
|
||||
# 1. a pid `pgrep` just listed can exit before the `ps -o comm=` lookup runs; on a gone pid `ps`
|
||||
# prints nothing, `comm` ends up empty, and an empty string matches none of the denied shell
|
||||
# names — so a pid that no longer exists was still counted.
|
||||
# 2. the denylist only knows the shells someone thought to name. `ssh`, `perl`, `python3`,
|
||||
# `ruby`, `tail` — anything else that carries the pattern in its own argv — was still
|
||||
# counted right along with the real daemon, and the ticket names `ssh` as a live route.
|
||||
# Both close with the same change: allowlist `comm = java` instead of denying shells. Measured on
|
||||
# the live daemon: `pid=30224 comm=java`. An empty comm (hole 1) is not `java` either, so it is
|
||||
# excluded for free — no separate "is this pid still alive" check needed.
|
||||
#
|
||||
# The objection, because it is real: an allowlist can UNDER-count. If fleetd ever stops being
|
||||
# launched as `java -jar ...` — a native image, a renamed launcher — `running_pid()` silently
|
||||
# returns nothing and `assert_single_daemon` stops noticing a second daemon at all. For a guard,
|
||||
# that false-negative direction is the worse one to be wrong in. This is not a new assumption,
|
||||
# though: `PATTERN='target/fleetd.jar'` two lines up already assumes the daemon is a jar, which
|
||||
# is only ever run by `java`. If that launch method changes, `PATTERN` stops matching anything
|
||||
# before this allowlist would ever get the chance to be wrong — the allowlist rides on the same
|
||||
# assumption that is already load-bearing, it does not add a new one. Whoever changes the launch
|
||||
# method needs to update both `PATTERN` and this allowlist together.
|
||||
running_pid() {
|
||||
local pid comm out=''
|
||||
for pid in $(pgrep -f "$PATTERN" 2>/dev/null || true); do
|
||||
comm="$(ps -o comm= -p "$pid" 2>/dev/null || true)"
|
||||
comm="${comm##*/}"
|
||||
comm="${comm#-}"
|
||||
# Allowlist, not a denylist of wrappers — see the CORRECTION 1 comment above. Anything that
|
||||
# is not literally `java` is excluded, including an empty comm from a pid that already exited.
|
||||
[ "$comm" = java ] || continue
|
||||
out="$out$pid"$'\n'
|
||||
done
|
||||
printf '%s' "$out"
|
||||
}
|
||||
|
||||
# fleetd #493 — three small, independently testable pieces of "never build into the path a
|
||||
# running process holds":
|
||||
@@ -507,8 +555,11 @@ assert_single_daemon() {
|
||||
die "more than one fleetd process is running after this restart (pids: $(printf '%s' "$pids" | tr '\n' ' ')).
|
||||
This is the exact failure a racing supervisor produces: the OLD jar was revived by its
|
||||
supervisor while this script started a NEW copy. Two daemons on one herdr session kill
|
||||
each other's members. Investigate with 'pgrep -f \"$PATTERN\"' and stop the wrong one by
|
||||
hand — do not assume either pid is the one you want."
|
||||
each other's members. Investigate with 'ps -eo pid,comm,args | grep -F \"$PATTERN\"' and
|
||||
check the COMM column of each hit yourself before acting — a bare 'pgrep -f \"$PATTERN\"'
|
||||
(fleetd #593) can match the very shell you type it into, not just the daemon, so it is not
|
||||
safe remediation advice on its own. Stop the wrong one by hand — do not assume either pid
|
||||
is the one you want."
|
||||
fi
|
||||
}
|
||||
|
||||
|
||||
@@ -433,6 +433,130 @@ test_assert_single_daemon_rejects_two_pids() {
|
||||
printf '%s' "$output" | grep -qF '4343' || fail "refusal message does not list the pids it found"
|
||||
}
|
||||
|
||||
# fleetd #593 instance 2 — `running_pid()` used to be a bare `pgrep -f "$PATTERN"`, which matches
|
||||
# ANY process whose full command line contains the pattern TEXT, including a shell that merely
|
||||
# embeds it as literal text rather than being the daemon. Measured live on this Mac: `pgrep -c`
|
||||
# (a one-call count) does not exist on BSD at all, and `bash -c "<single command>"` execs in place
|
||||
# so no parent shell survives to hold the pattern — which is exactly why the defect did not
|
||||
# reproduce from a plain script and needs a wrapper shaped like this instead. A `sh -c '...; ...'`
|
||||
# with MORE THAN ONE statement does not get that exec-in-place treatment: the shell forks a child
|
||||
# for the second statement and stays alive itself, holding the whole `-c` string — pattern text
|
||||
# included — in its own `ps -o args`, for as long as it runs. That is the same shape an
|
||||
# `ssh host "…; …"` wrapper or a hand-typed pipeline leaves behind. Before the fix this test would
|
||||
# have found the wrapper's pid in running_pid()'s output; it must not.
|
||||
test_running_pid_excludes_self_matching_wrapper_shell() {
|
||||
local before after wrapper_pid
|
||||
before="$(running_pid)"
|
||||
sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' &
|
||||
wrapper_pid=$!
|
||||
sleep 0.3
|
||||
after="$(running_pid)"
|
||||
kill "$wrapper_pid" 2>/dev/null || true
|
||||
wait "$wrapper_pid" 2>/dev/null || true
|
||||
[ "$after" = "$before" ] \
|
||||
|| fail "running_pid() counted a self-matching wrapper shell (pid $wrapper_pid, holding the pattern as literal text in its own argv, not the daemon): before=[$before] after=[$after]"
|
||||
}
|
||||
|
||||
# fleetd #593 CORRECTION 1 — the round-1 version of this test gave its standin an argv[0]
|
||||
# containing the pattern text (via `exec -a`) and left `comm` as whatever that override produced,
|
||||
# which was never `java`. That was fine for a denylist-of-shells filter, but the allowlist below
|
||||
# now requires `comm = java` specifically, so the standin here must actually carry that comm, not
|
||||
# just avoid being a shell. `exec -a java` overrides argv[0] to `java` while the process itself
|
||||
# stays a genuine, harmless `sh`; combining it with the same non-exec'ing multi-statement shape
|
||||
# the wrapper-shell test above uses keeps the pattern text in the process's own `ps -o args` for
|
||||
# as long as it runs. Measured live on this Mac (BSD/macOS: `ps -o comm=` here reflects argv[0]):
|
||||
# `comm=java`, `args` contains the pattern, `pgrep -f "$PATTERN"` finds it. Copying a real system
|
||||
# binary into a scratch path and executing it from there was tried first, for a more literal
|
||||
# stand-in daemon, and the OS killed it outright (SIGKILL, exit 137 — almost certainly a
|
||||
# code-signing check on a relocated binary); `exec -a` needs no binary of its own and nothing
|
||||
# under a scratch directory, and it is the technique CORRECTION 1 names as the right one.
|
||||
#
|
||||
# This is the one live-process test in this file whose result could differ on Linux: Linux sets
|
||||
# `comm` from the actually-executed binary's own path, not from `exec -a`'s argv[0] override (BSD
|
||||
# ties `comm` to argv[0], which is what makes this technique work here) — so on Linux this
|
||||
# specific fixture might report `comm=sh`, not `comm=java`, even though the REAL daemon (a literal
|
||||
# `java -jar target/fleetd.jar` process, never fabricated) is unaffected either way. I could not
|
||||
# verify this fixture's behavior on Linux, so test_running_pid_counts_a_pid_whose_comm_is_java
|
||||
# below backstops the same claim (the allowlist admits a pid whose comm is `java`) with a stubbed
|
||||
# `ps`, which is identical bash on every platform and carries no such platform question.
|
||||
test_running_pid_finds_a_real_java_named_second_process() {
|
||||
local before after standin_pid
|
||||
before="$(running_pid)"
|
||||
( exec -a java sh -c 'echo "target/fleetd.jar" >/dev/null; sleep 20' ) &
|
||||
standin_pid=$!
|
||||
sleep 0.3
|
||||
after="$(running_pid)"
|
||||
kill "$standin_pid" 2>/dev/null || true
|
||||
wait "$standin_pid" 2>/dev/null || true
|
||||
printf '%s\n' "$after" | grep -qxF "$standin_pid" \
|
||||
|| fail "running_pid() did not find a real second process (pid $standin_pid, comm forced to 'java' via exec -a) whose own argv holds the pattern: before=[$before] after=[$after]"
|
||||
}
|
||||
|
||||
# fleetd #593 CORRECTION 1, hole 2 — the round-1 filter denied known shell names (sh/bash/zsh/
|
||||
# dash/ksh) and counted everything else. `ssh`, `perl`, `python3`, `ruby`, `tail` — anything not on
|
||||
# that list, carrying the pattern in its own argv — was still counted right alongside the real
|
||||
# daemon, and the ticket names `ssh` as a live route. Stubbing `pgrep`/`ps` (rather than spawning a
|
||||
# real perl/ssh process) pins the exact discriminator this correction is about — comm, not the
|
||||
# caller's shape — deterministically on every platform, with no dependency on perl/python3/ruby
|
||||
# being installed in whatever environment runs this suite, and no dependency on how a given OS
|
||||
# derives `comm` for a fabricated process (see the comment above
|
||||
# test_running_pid_finds_a_real_java_named_second_process for why that matters here).
|
||||
test_running_pid_drops_a_pid_whose_comm_is_not_java() {
|
||||
pgrep() { printf '4242\n'; }
|
||||
ps() { printf 'perl\n'; }
|
||||
local found
|
||||
found="$(running_pid)"
|
||||
unset -f pgrep ps
|
||||
[ -z "$found" ] \
|
||||
|| fail "running_pid() counted pid 4242 whose comm is 'perl', not 'java' — denying known shell names does not exclude a non-shell wrapper such as ssh or perl (fleetd #593 CORRECTION 1): found=[$found]"
|
||||
}
|
||||
|
||||
# fleetd #593 CORRECTION 1, hole 1 — pgrep can list a pid that exits before the following
|
||||
# `ps -o comm=` lookup runs; on a gone pid `ps` prints nothing, so `comm` comes back empty. Under
|
||||
# the round-1 denylist an empty string matched none of the denied shell names, so the dead pid was
|
||||
# still counted — the exact false-positive shape the ticket exists to remove, just rarer. The
|
||||
# allowlist fixes this for free: an empty comm is not `java` either.
|
||||
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup() {
|
||||
pgrep() { printf '4242\n'; }
|
||||
ps() { :; } # a pid that no longer exists: the real `ps -p <gone>` prints nothing and this mirrors that
|
||||
local found
|
||||
found="$(running_pid)"
|
||||
unset -f pgrep ps
|
||||
[ -z "$found" ] \
|
||||
|| fail "running_pid() counted pid 4242 whose comm lookup came back empty (the pid had already exited before the lookup ran) — an empty comm must not pass the allowlist (fleetd #593 CORRECTION 1): found=[$found]"
|
||||
}
|
||||
|
||||
# The positive backstop for both stubbed tests above, and for
|
||||
# test_running_pid_finds_a_real_java_named_second_process on whatever platform that live fixture
|
||||
# does not itself carry comm=java: the allowlist must still ADMIT the one comm value the real
|
||||
# daemon actually has. Measured on the real, currently-running daemon on this Mac: `comm=java`.
|
||||
test_running_pid_counts_a_pid_whose_comm_is_java() {
|
||||
pgrep() { printf '4242\n'; }
|
||||
ps() { printf 'java\n'; }
|
||||
local found
|
||||
found="$(running_pid)"
|
||||
unset -f pgrep ps
|
||||
printf '%s\n' "$found" | grep -qxF '4242' \
|
||||
|| fail "running_pid() did not count pid 4242 whose comm is 'java' — the daemon's own name must pass the allowlist: found=[$found]"
|
||||
}
|
||||
|
||||
# fleetd #593 instance 3 — assert_single_daemon's refusal message used to tell the operator to
|
||||
# "Investigate with 'pgrep -f \"\$PATTERN\"'", which — typed by hand or over ssh — is precisely the
|
||||
# self-matching invocation instance 2 above fixes. A source-text check, the same technique
|
||||
# test_no_error_lines_message_gated_by_drain_state uses: this is prose inside a die() call, never
|
||||
# reached by sourcing (the SOURCED guard stops before the main flow, and this text only prints
|
||||
# from inside a call assert_single_daemon makes when it is already refusing).
|
||||
test_die_message_does_not_recommend_bare_pgrep_as_remediation() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" block bad
|
||||
block="$(grep -A6 -F 'racing supervisor produces' "$src" || true)"
|
||||
[ -n "$block" ] || fail "could not find the assert_single_daemon refusal message in redeploy-fleetd.sh"
|
||||
bad="$(printf '%s' "$block" | grep -F "Investigate with 'pgrep -f" || true)"
|
||||
[ -z "$bad" ] \
|
||||
|| fail "assert_single_daemon's die message still hands the operator a bare 'pgrep -f \"\$PATTERN\"' as remediation (fleetd #593) — that is exactly the self-matching invocation"
|
||||
printf '%s' "$block" | grep -qF 'fleetd #593' \
|
||||
|| fail "assert_single_daemon's die message does not say in words that a pattern can match the caller (fleetd #593)"
|
||||
}
|
||||
|
||||
# fleetd #511 — jar_id()'s no-argument default was unpinned by any test: nothing proved it reports
|
||||
# $JAR (the live path) rather than $JAR_STAGED. Both halves matter, so this pins both: the bare call
|
||||
# must hash the live jar, and an explicit path argument must hash THAT file, not fall back to $JAR.
|
||||
@@ -1894,6 +2018,12 @@ test_require_drivable_supervisor_accepts_known_kinds
|
||||
test_count_daemon_pids
|
||||
test_assert_single_daemon_accepts_one_pid
|
||||
test_assert_single_daemon_rejects_two_pids
|
||||
test_running_pid_excludes_self_matching_wrapper_shell
|
||||
test_running_pid_finds_a_real_java_named_second_process
|
||||
test_running_pid_drops_a_pid_whose_comm_is_not_java
|
||||
test_running_pid_drops_a_pid_that_exited_before_the_comm_lookup
|
||||
test_running_pid_counts_a_pid_whose_comm_is_java
|
||||
test_die_message_does_not_recommend_bare_pgrep_as_remediation
|
||||
test_jar_id_defaults_to_live_and_reports_explicit_path
|
||||
test_hash256_computes_a_real_sha256
|
||||
test_jar_id_reports_absent_for_missing_file
|
||||
|
||||
Reference in New Issue
Block a user