Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 6cccd458d4 #589 Group 3: wiring-test the 5 sites below line 500 in Fleetd.main()
CI / shell-tests (pull_request) Successful in 11s
CI / contract (pull_request) Successful in 1m26s
CI / build (pull_request) Successful in 2m11s
Extracts the inline lambdas/method references at the 5 assigned wiring
sites into named package-private factories on Fleetd, following the
FleetdLoopHealthSourceWiringTest pattern from #584:

- turnRegistrar(CompletionResolver) — was completion::register (Injector)
- healthFailTarget(MessageService) — was messages::abandon (FleetHealthMonitor)
- releaseCleanup(MessageService, ReplyInbox, PrimaryRegistry) — was the
  inline sessions.onRelease(detail -> {...}) cleanup lambda
- replyInboxOpener() — was AmqpReplyInbox::open passed to selectReplyInbox
- leadMailboxOpener() — was LeadMailbox::open passed to openLeadMailbox

Each factory has a new runtime test (not source-text) that drives real
collaborators through public APIs: MessageService.poll(ticket).phase(),
InMemoryReplyInbox.peek(), PrimaryRegistry.nudgeTargetFor(), and the
opener tests connect to a guaranteed-closed local port to prove a real
network attempt vs. an inert stub.

releaseCleanup was done first per the brief: MessageService.abandon's
javadoc documents that losing this cleanup leaves a torn-down worker's
rendezvous waiter open forever.

Tests: 1789 -> 1794 (+5), 0 failures, 0 errors. mvn -q -o test exit 0,
no BUILD FAILURE, no piped exit status. Each new test verified RED on
the inert form named in the ticket, and GREEN after reformatting the
call across lines and extracting the argument into a local/factory.
2026-09-19 15:24:11 +07:00
12 changed files with 493 additions and 652 deletions
+137 -161
View File
@@ -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).
// 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);
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
// 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(claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
claudeProfiles, cfg, config));
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));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(openCodeLauncher(router.memberAgents(), router.memberSpaces(),
opencodeProfiles, cfg, config, forwardingExhaustionSink));
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));
}
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.
// 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);
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);
// 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,13 +464,11 @@ 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.
// 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);
exhaustionSinkRef.set(exhaustionSink);
// 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
@@ -503,21 +504,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.
@@ -587,10 +593,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)) {
@@ -610,27 +618,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.
@@ -1011,120 +1003,6 @@ 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
@@ -1182,6 +1060,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).
@@ -1,84 +0,0 @@
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");
}
}
@@ -1,86 +0,0 @@
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"));
}
}
@@ -1,62 +0,0 @@
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");
}
}
@@ -1,110 +0,0 @@
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,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");
}
}
@@ -1,72 +0,0 @@
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"));
}
}
@@ -1,77 +0,0 @@
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");
}
}