fleetd #612 Unit A: extract Fleetd.main's boot composition into FleetdAssembly/FleetdRuntime
CI / shell-tests (pull_request) Failing after 12s
CI / build (pull_request) Failing after 1m29s
CI / contract (pull_request) Successful in 1m28s

Fleetd.main kept config loading, startup reports and validation. Everything from
the herdr socket connect onward moved verbatim, same order, into
FleetdAssembly.assembleAndStart(AssemblyInputs, ResourcePorts), which returns a
FleetdRuntime owning the real objects (package-private accessors, never a copy)
and their single ordered close(). ResourcePorts/SystemResourcePorts abstract every
boot-time side effect (env, herdr connect, broker openers, clocks, schedulers,
shutdown-hook registration, HTTP start) with no inert production variant, per the
architect proposal on the ticket.

sleepHerdrPoll widened from private to package-private so FleetdAssembly can pass
a method reference to it; no other signature changed.

FleetdAssemblyLifecycleTest drives the real assembly with FakeHerdr, a temp
FleetConfig and a fake ResourcePorts recording a start/close ledger, asserting it
against the order recorded from the pre-move main() and shutdown hook, and proving
every resource the ledger can observe (herdr client, three schedulers, the AMQP
reply inbox) closes via FleetdRuntime.close(). FakeHerdr gained a closed flag for
this.
This commit is contained in:
Dai Ha
2026-09-22 10:52:35 +07:00
parent 17127efb88
commit 3f7bc3815e
8 changed files with 1108 additions and 592 deletions
@@ -0,0 +1,22 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
/**
* fleetd #612 Unit A: everything {@code Fleetd.main} has ready once config is loaded, reported
* and validated — the exact point {@code main} used to keep going straight into socket and broker
* work. Building this (and handing it to {@link FleetdAssembly#assembleAndStart}) is the seam a
* test now has to drive the real boot composition without being {@code main} itself.
*
* @param cfg the boot-time {@link FleetConfig} snapshot every one-time wiring decision reads —
* see {@code Fleetd.main}'s own comment on why this must never be swapped for a live
* reference once loaded
* @param config the live {@link ConfigRef} the hot-reloadable paths read per use
* @param guard the same {@link SubscriptionGuard} {@code Fleetd.main} already used to assert the
* launching environment is clean, reused rather than rebuilt so the assembled
* launchers see the identical instance {@code main} already validated with
*/
record AssemblyInputs(FleetConfig cfg, ConfigRef config, SubscriptionGuard guard) {
}
+10 -592
View File
@@ -195,595 +195,10 @@ public final class Fleetd {
// above as ConfigRef's extraValidation, so a reload refuses what this line refuses at startup.
assertChartersNameOnlyRegisteredTools(cfg);
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
? Path.of(cfg.herdrSocket())
: UnixSocketHerdrClient.defaultSocketPath();
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
UnixSocketHerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()
? UnixSocketHerdrClient.connect(Path.of(cfg.memberHerdrSocket()), new com.fasterxml.jackson.databind.ObjectMapper())
: herdr;
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
target -> leadsRef.get().get().containsKey(target));
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
Map<String, FleetConfig.Profile> claudeProfiles = new LinkedHashMap<>();
Map<String, FleetConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
cfg.profiles().forEach((name, w) -> {
if (w.isOpenCode()) {
opencodeProfiles.put(name, w);
} else {
claudeProfiles.put(name, w);
}
});
List<HerdrPeerLauncher> adapters = new ArrayList<>();
// fleetd #175: the daemon's real ExhaustionSink can only be built once `sessions` exists
// (below), but `sessions` needs `workers`, which needs the adapters built right here — a
// genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now,
// pointed at the real one once it exists.
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
// fleetd #234, round 4: routed through the shared ExhaustionSink.forwardingTo factory
// rather than written inline here — not because a lambda at this call site is unsafe
// anymore (it is not: the 3-arg overload is now the interface's single abstract method, so
// 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);
// 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));
}
if (!opencodeProfiles.isEmpty()) {
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
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
// here, at startup, and a config reload only changes it for a daemon restart.
// fleetd #466: escalating, not flat — a credential that keeps reporting exhaustion (e.g. a
// weekly subscription limit, which would otherwise be retried on every ~30-minute cooldown,
// about 336 times across the week) backs off further each consecutive time, capped at
// BackendQuarantine.DEFAULT_MAX_COOLDOWN_MULTIPLE x the base cooldown. See BackendQuarantine's
// class doc for the mechanism, why this never fires on cooling-off (a separate, unescalated
// mechanism — BackendOutagePolicy below), and the reset.
BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
// fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the
// launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in
// below (written on a classified backend error). A SEPARATE, shorter-lived mechanism from
// `quarantine` — see BackendOutagePolicy's class doc — never merged with it.
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(System::nanoTime);
PeerLauncher workers = new CompositePeerLauncher(
adapters,
cfg.effectiveDefaultProfile(),
config,
profileName -> liveCountRef.get().apply(profileName),
quarantine,
outagePolicy);
// fleetd #422 follow-up: say which of the three model-gate states the daemon booted into —
// no models: block at all, a block armed with nothing off, or a block with N off — the same
// way exhaustedPatternCoverageLine/errorPatternCoverageLine report CB-578 stage A/fleetd
// #201 Unit 5 coverage just below. Read from workers.modelGateState() (never a separate
// config.get().models() here) so this line and fleet_profiles' modelGateArmed can never
// disagree about what CompositePeerLauncher's spawn gate actually enforces.
log.info("model gate (fleetd #422): {}", modelGateCoverageLine(workers.modelGateState()));
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
// the first thing that actually talks to herdr, so without this wait a boot-order race
// would crash the daemon into a restart loop. Wait, then degrade rather than die: serving
// with /healthz reporting "degraded" is strictly more useful than exiting.
HerdrAwaitOutcome herdrOutcome = awaitHerdr(herdr, System::nanoTime, Fleetd::sleepHerdrPoll);
boolean herdrUp = logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
if (herdrUp) {
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
// with the previous process — reap those leaked orphans now, before we start serving.
workers.reapOrphanWorkers();
}
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
// CB-303 part 2: context cap is opt-in and disabled (0) when absent/null.
int contextCap = 0;
if (cfg.lifecycle() != null && cfg.lifecycle().contextCap() != null
&& cfg.lifecycle().contextCap() > 0) {
contextCap = cfg.lifecycle().contextCap();
}
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
SessionManager sessions = new SessionManager(workers,
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
System::nanoTime, contextCap, clearAfterTurn);
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
// Idle-sleep guard: hold an OS-level assertion against idle sleep while at least one
// member is live, so an unattended host does not idle-sleep out from under a member's
// long turn (see FleetConfig.IdleSleepGuard / dev.ltms.fleet.power.IdleSleepGuard for the
// measurement that motivated this). Opt-out via idleSleepGuard.enabled: false; on by
// default. Hangs off SessionManager's own onAcquire/onRelease hooks (CB-520/CB-516,
// previously wired only to the reply inbox) and SessionManager#size() — the exact registry
// fleet_list's live/capacity numbers are themselves computed from — rather than tracking
// members a second way. No-op (never constructed) off macOS or when idleSleepGuard.enabled
// is explicitly false; the mechanism itself is additionally a no-op if 'caffeinate' cannot
// be started, so this can never fail a spawn, a release, or startup.
boolean idleSleepGuardEnabled = cfg.idleSleepGuard() == null || cfg.idleSleepGuard().isEnabled();
final IdleSleepGuard idleSleepGuard;
if (idleSleepGuardEnabled) {
idleSleepGuard = new IdleSleepGuard(new CaffeinateSleepAssertionMechanism(), sessions::size);
sessions.onAcquire(_ -> idleSleepGuard.recheck());
sessions.onRelease(_ -> idleSleepGuard.recheck());
} else {
idleSleepGuard = null;
}
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
final SessionReaper reaper;
if (cfg.lifecycle() != null
&& cfg.lifecycle().idleTtlSeconds() != null
&& cfg.lifecycle().idleTtlSeconds() > 0) {
reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds());
reaper.start();
} else {
reaper = null;
}
// CB-530: every pane the config names as a lead, merged from `leaders:` and the legacy
// singular pin. PrimaryRegistry below still tracks ONE terminal — it addresses the push
// loop's nudges, which need a single destination — so it keeps the legacy pin.
Map<String, String> leadTerminals = cfg.leaderTerminals();
if (leadTerminals.size() > 1) {
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
}
// CB-531: on top of the legacy primary.terminal pin, discover leads by the tab labels the
// operator writes. CB-557 moved the settings onto the lead they describe, so scanning is on
// whenever a `fleet.leaders:` entry exists — with no leads configured the supplier is a
// constant and never touches herdr, exactly as a missing `leadScan:` block used to behave.
// CB-579: each lead now names its own exact `tab:` label, so one scanner discovers every
// configured lead regardless of how differently their tabs are labelled — the old
// single-shared-tabPrefix limitation (and its warning) is gone.
final Supplier<Map<String, String>> leads;
var leaders = cfg.fleet().leaders();
if (!leaders.isEmpty()) {
Map<String, String> tabToName = new LinkedHashMap<>();
leaders.forEach((name, leader) -> {
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
tabToName.put(leader.tab(), name);
}
});
// The lead and the members now share ONE workspace (the operator asked for a single
// "session" with many tabs), so no workspace can be excluded — the lead lives in the
// members' space by design. A lead is told from a member by its exact tab label alone:
// a lead carries its configured `tab`, a member its `worker: {profile} #{n}` template,
// and the two never collide. (The scanner still supports an exclusion set for a split
// layout; the fleet's policy here is simply not to use one.)
//
// One shared rescan cadence: still taken from the first entry, as before — it is an
// operational cadence, not identity, so there is no correctness reason to give every
// lead its own scanner.
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), System::nanoTime);
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
tabToName.keySet(), scanIntervalSeconds);
} else {
leads = () -> leadTerminals;
}
leadsRef.set(leads);
// CB-558: start any declared lead that is not already running. After the scanner is built,
// because both read the same tab labels and the ordering makes that dependency visible; and
// only when herdr answered, because the launcher's whole safety property is that it can
// count live leads first — it must never guess and risk a second orchestrator.
if (herdrUp && !leaders.isEmpty()) {
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
if (launched > 0) {
log.info("lead auto-launch: {} lead(s) started", launched);
}
}
// CB-548: config-declared architect slots. Config supplies only the stable name → profile
// map; the terminal → slot binding is owned by the registry and is empty at startup, so no
// pane resolves to an architect until the later spawn lifecycle binds one. The registry is
// what CallerResolver resolves against and what that lifecycle will read profiles from;
// nothing here spawns a slot.
// fleetd #424: MemberRegistry.live re-reads fleet.architects through `config` on every
// reserve/requireSlotFor call, so a reload that removes or adds an architect slot governs
// the next spawn with no restart — the frozen `new MemberRegistry(cfg.fleet())` this used
// to be let a "revoked" slot keep granting new architect spawns forever.
MemberRegistry members = MemberRegistry.live(() -> config.get().fleet());
sessions.setMemberLifecycle(members);
if (!members.slots().isEmpty()) {
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
+ "spawn lifecycle binds a live terminal to it)",
members.slots().size(), members.slots().keySet());
}
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then.
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
Rendezvous rendezvous = new Rendezvous();
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
// fleetd #446: read live off `config` per lookup, cached by profile name — see
// LiveExhaustedPatterns's class doc for why this replaces the old compiled-once-at-startup
// map. A profile with no exhaustedPattern simply returns null here, so its workers keep
// today's completion-fallback behaviour unchanged.
// 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.
Set<String> exhaustedConfiguredAtStartup = cfg.profiles().entrySet().stream()
.filter(e -> e.getValue().hasExhaustedPattern())
.map(Map.Entry::getKey)
.collect(Collectors.toCollection(LinkedHashSet::new));
log.info("backend-exhausted classification (CB-578 stage A): {}",
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedConfiguredAtStartup));
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
// rather than handing it back as a real answer. Still compiled once at startup, keyed by
// profile name — unlike exhaustedPattern above (fleetd #446), errorPattern was left DEFERRED
// on purpose: the ticket that made exhaustedPattern hot scoped errorPattern/cooling-off out
// explicitly. A profile with no configured errorPattern is simply absent here, so
// CompletionResolver falls back to its built-in narrow {@code (?i)\bAPI Error\s*:}
// compatibility pattern for that profile's targets (BackendErrorPatternLookup#legacy's
// contract — see backendErrorPatterns below).
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasErrorPattern()) {
errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern()));
}
});
// fleetd #248: extracted to a static factory (see backendErrorPatternLookup below) so a
// test can prove main() actually PASSES this into CompletionResolver, not only that the
// lookup itself behaves correctly — the exact gap fleetd #248 exists to close.
BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,
errorPatternsByProfile);
log.info("backend-error classification (fleetd #201 Unit 5): {}",
errorPatternCoverageLine(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
// fleetd #446 criterion 3: the backend text that triggered the most recent quarantine of
// each credential, so fleet_profiles can report WHY a limit was hit, not only that it was.
// Keyed by credential id — the same key BackendQuarantine's own remainingSeconds uses —
// and written at the one call site that actually quarantines (inside exhaustionSink below),
// so a reason can never be reported for a quarantine that never happened. Bounded the same
// way BackendQuarantine's own internal map is documented to be: by the number of distinct
// credentials ever exhausted, not by how often the config is edited.
Map<String, String> quarantineReasonByCredential = new ConcurrentHashMap<>();
// fleetd #446 follow-up (round 3): extracted into exhaustionSink(...) below — see that
// method's javadoc for the full fleetd #175/#234/#446 history this used to carry inline —
// so a dedicated test can drive the exact ExhaustionSink main() builds, not a hand-rebuilt
// copy of its shape.
// 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);
// 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
// holder set once `pushLoop` exists, read lazily from inside the lambda built here.
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>();
// fleetd #248: extracted to a static factory (see backendErrorSink below), public rather
// than package-private like the other two factories here, so
// dev.ltms.fleet.inject.BackendOutageFlowTest can exercise the REAL production sink
// directly instead of a hand-mirrored copy of this lambda — that copy was precisely the
// gap fleetd #248 exists to close (see that test's class doc for the history).
BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),
outagePolicy, pushLoopRef::get);
AgentControl agents = router.memberAgents();
// Both fleetd#201 Unit 5 (backend-error patterns + sink) and fleetd#241 (the worktree/branch
// lookup the fallback report names) land on this one call. The full constructor takes both,
// so neither feature is dropped; nowNanos must be passed explicitly to reach it.
//
// fleetd #248: every argument built specifically for this call (backendErrorPatterns and
// backendErrorSink above, and the worktree/branch lookup right here) now comes from a
// static factory tested on its own; FleetdCompletionResolverWiringTest source-asserts that
// THIS call actually passes them, which is the coverage that was missing before.
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,
worktreeBranchLookup(sessions::roster));
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
// fleetd #561: extracted to a static factory (see turnListener below) — the anonymous class
// this replaced had two bare, unguarded statements per callback, and nothing enforced that
// the completion half went first beyond call order in the source.
TurnListener turnListener = turnListener(completion, sessions);
Predicate<String> deliverable = deliverableTo(presence, leads);
// 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, 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.
// 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.
// 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.
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs —
// identity comes from `leaders:`/`leadScan:`, and reply nudges now follow the delegating
// lead. Say so once at startup rather than leaving a redundant pin to look load-bearing.
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
+ "from leaders:/leadScan:, and reply nudges follow the lead that delegated. "
+ "It still works, and is still the fallback nudge destination when a restart "
+ "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid.");
}
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-push-").unstarted(r));
// CB-502: the registry is built before the service and the push loop so send/reply outcomes
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
// CB-512: the push loop takes it too, so nudge outcomes (delivered|exhausted) are counted.
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
// above at the real push loop, now that it exists.
pushLoopRef.set(pushLoop);
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
// It has its own single-thread scheduler and holds its own scheduler shutdown via close().
final LeadHeartbeatLoop heartbeat;
var heartbeatScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
if (cfg.leadHeartbeat() != null) {
var hb = cfg.leadHeartbeat();
// fleetd #609: own LeadContextGauge instance for the heartbeat loop — separate from the
// one FleetMcp builds internally for fleet_list's context row. Each caches independently
// (keyed by configDir+sessionId), so this costs at most one extra bounded tail read per
// TTL window, never a shared-mutable-state hazard between the two callers.
var leadContextGauge = new LeadContextGauge();
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
pushLoop, heartbeatScheduler, System::nanoTime,
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
metrics,
leadContextSource(leadContextGauge, router.leadAgents(), leads,
leadConfigDirLookup(() -> config.get().profiles(), leaders)),
Boolean.TRUE.equals(hb.contextHighNudge()));
heartbeat.start();
} else {
heartbeat = null;
heartbeatScheduler.shutdownNow();
}
// fleetd #480: lead rollover. Opt-in; absent `leadRollover:` this is never constructed, so
// an upgraded daemon cannot silently acquire the ability to clear the lead's own pane.
// Unlike heartbeat above, this has no recurring scheduler of its own — nothing but an
// explicit confirm() call (wired to an MCP tool by a later ticket; nothing calls it yet)
// that passes every gate can ever schedule a roll. It does not take primaryRegistry: the
// lead terminal to roll comes from the caller of open()/confirm() (resolved by the MCP
// layer from the connection, the same way auth/CallerResolver#resolve builds a
// Principal.leader(...)), never from a single-slot lookup — see LeadRollover's class
// javadoc, fleetd #480 correction 2.
LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config, leads);
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
pushLoop, metrics);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
final FleetHealthMonitor healthMonitor;
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-health-").unstarted(r));
if (cfg.health() != null && cfg.health().isEnabled()) {
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
// through the same idempotent target-wide operation CB-516 already uses on release.
// 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(), healthFailTarget(messages));
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
log.warn("fleet health: {} (no notification sink configured)", coverage);
} else {
log.info("fleet health: {}", coverage);
}
healthMonitor.start();
} else {
healthMonitor = null;
healthScheduler.shutdownNow();
}
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
sessions.onAcquire(replyInbox::own);
// 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.
// 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.
// CB-185: a caller's pane can live on either daemon (a lead's on the lead daemon, a
// member's on the member daemon) — search both, lead first. Collapses to one scan when
// memberHerdrSocket is unset (herdr == memberHerdr).
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
final CallerResolver callers;
if (cfg.auth().tokenMode()) {
String token = System.getenv(cfg.auth().tokenEnv());
if (token == null || token.isBlank()) {
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
+ " is unset or empty — export it before starting fleetd");
}
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
log.info("auth: token mode (bearer required for non-worker callers, env {})",
cfg.auth().tokenEnv());
} else {
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
// fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
// closure in only one of the two places), exactly the shape #284 was.
FleetMcp.QuarantineSource quarantineSource = quarantineSource(config, quarantine,
liveExhaustedPatterns, quarantineReasonByCredential);
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp.LoopHealthSource loopHealth = loopHealthSource(poller, reaper);
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
healthCoverageSource(config),
loopHealth,
quarantineSource,
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
// leadMailbox above), so peers follows the same rule rather than half hot-reloading.
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers(),
// fleetd #480 Unit C: the executor behind fleet_handover — null whenever
// leadRollover: is not configured (see the leadRollover local above).
leadRollover);
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
// coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is
// created and no thread runs. It reads the SAME live lead supplier the injector's
// deliverability gate does, so a lead found by the tab scan after startup is reachable
// without a restart.
final LeadCoordLoop leadCoordLoop;
final ScheduledExecutorService leadCoordSchedulerRef;
if (leadMailbox != null) {
var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r));
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
LEAD_COORD_INTERVAL_MS);
leadCoordLoop.start();
leadCoordSchedulerRef = leadCoordScheduler;
} else {
leadCoordLoop = null;
leadCoordSchedulerRef = null;
}
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
final ConfigWatcher configWatcher;
if (cfg.configReload() != null && cfg.configReload().isEnabled()) {
configWatcher = new ConfigWatcher(config, cfg.configReload().intervalSeconds());
configWatcher.start();
} else {
configWatcher = null;
}
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
// open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr
// last. This replaces the earlier independent hooks that could race and close herdr early.
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
poller.stop();
messages.close();
pushLoop.close();
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
if (leadCoordLoop != null) leadCoordLoop.close(); // CB-637: stop delivering peer-lead messages
if (leadCoordSchedulerRef != null) leadCoordSchedulerRef.shutdownNow();
if (healthMonitor != null) healthMonitor.stop();
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
mcp.close();
if (reaper != null) reaper.stop();
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
// drained every session (and each release already drove the live count to 0, which
// releases the guard's assertion on its own) — this is the backstop for a drain that was
// itself interrupted or threw, so no caffeinate child ever outlives the daemon.
if (idleSleepGuard != null) idleSleepGuard.close();
// Release the broker connection last among message resources (no-op for the in-memory inbox).
if (replyInbox instanceof AutoCloseable closeable) {
try {
closeable.close();
} catch (Exception e) {
log.debug("reply inbox close: {}", e.toString());
}
}
// CB-637: the coordination connection goes with it — after the loop that reads it has
// stopped, so no tick can be mid-ack against a closed channel.
if (leadMailbox != null) {
try {
leadMailbox.close();
} catch (Exception e) {
log.debug("lead mailbox close: {}", e.toString());
}
}
router.close();
}));
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
// fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials —
// same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above.
// fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above —
// GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles.
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable,
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
quarantineSource, outageSource, loopHealth).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("fleetd listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
// fleetd #612 Unit A: everything from here on used to run inline in this method. It now
// lives in FleetdAssembly.assembleAndStart, built against a real ResourcePorts — see that
// class's javadoc for the full boot-order contract this preserves exactly.
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ResourcePorts.system());
}
/**
@@ -2286,13 +1701,16 @@ public final class Fleetd {
record HerdrAwaitOutcome(HerdrWaitResult result, long elapsedNanos) {}
/**
* The real per-poll wait {@link #main} passes to {@link #awaitHerdr}: sleep
* {@link #HERDR_WAIT_POLL_MILLIS}, and on interruption re-set the thread's interrupt flag
* The real per-poll wait {@link FleetdAssembly#assembleAndStart} passes to {@link #awaitHerdr}:
* sleep {@link #HERDR_WAIT_POLL_MILLIS}, and on interruption re-set the thread's interrupt flag
* rather than throwing — {@link #awaitHerdr} detects an interruption by checking {@link
* Thread#isInterrupted()} right after this returns, so a poller that swallowed the flag
* instead of restoring it would make that check silently miss the interruption.
*
* <p>Package-private (fleetd #612 Unit A), not {@code private}: {@code FleetdAssembly}, a
* different class in this package, needs a {@code Runnable} reference to this exact method.
*/
private static void sleepHerdrPoll() {
static void sleepHerdrPoll() {
try {
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
} catch (InterruptedException ie) {
@@ -0,0 +1,515 @@
package dev.ltms.fleet;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.ConfigWatcher;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.health.FleetHealthMonitor;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.herdr.LeadTabScanner;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
import dev.ltms.fleet.inject.BackendErrorPatternLookup;
import dev.ltms.fleet.inject.BackendErrorSink;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.lead.LeadRollover;
import dev.ltms.fleet.mcp.ConnectionIdentity;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.mcp.LsofPeerPidLookup;
import dev.ltms.fleet.mcp.LsofProcessCwdLookup;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
import dev.ltms.fleet.power.IdleSleepGuard;
import dev.ltms.fleet.rest.FleetApp;
import dev.ltms.fleet.session.GitWorktrees;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.SessionReaper;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.Predicate;
import java.util.function.Supplier;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
/**
* fleetd #612 Unit A: the real boot assembly, extracted out of {@code Fleetd.main} so a test can
* drive it directly. {@link #assembleAndStart} is <em>the same statements {@code main} used to run
* inline</em>, in the same order, against a real {@link ResourcePorts} in production and a fake one
* in a test — see {@code FleetdAssemblyLifecycleTest}. {@code Fleetd.main} keeps config loading,
* the startup reports and validation; everything from the herdr socket onward moved here.
*
* <p><strong>Construction and start order is preserved exactly, on purpose.</strong> This is not
* rebuilt into "construct everything, then start everything" — that would change boot timing. The
* order recorded before any code moved (see the ticket and {@code FleetdAssemblyLifecycleTest}):
* {@code SessionReaper} starts first (if {@code lifecycle.idleTtlSeconds} is configured), then
* {@link StatusPoller}, then the optional {@link LeadHeartbeatLoop} and {@link FleetHealthMonitor},
* then the optional {@link LeadCoordLoop} and {@link ConfigWatcher}, and the HTTP server starts
* last of all. The close order (see {@link FleetdRuntime#close()}) is the mirror the original
* shutdown hook always used.
*
* <p><strong>One statement could not move without reordering startup.</strong> {@code Fleetd.main}
* registered its shutdown hook <em>before</em> building the Javalin {@code FleetApp} — the hook
* itself never touched {@code app} (it still doesn't; see {@link FleetdRuntime#close()}), but the
* hook needs a {@link FleetdRuntime} to close over, and the ticket asks that runtime to also own
* {@code FleetApp}. Building the runtime before the app exists and mutating it afterward (via
* {@link FleetdRuntime#attachApp}) preserves the exact original order — hook registered, then app
* built, then HTTP started — without moving the app's construction earlier or the hook's
* registration later. That is the one seam this ticket did not get to pin any other way.
*
* <p><strong>No inert variant.</strong> Deliberately, there is no overload of this method that
* accepts a smaller/optional {@link ResourcePorts} or defaults one internally. A future edit that
* wants to skip {@code FleetdAssembly} entirely and build its own graph is still possible — no
* static analysis stops that — but it cannot do so by quietly swapping this call for an inert
* substitute that still compiles, because none exists.
*/
final class FleetdAssembly {
private static final Logger log = LoggerFactory.getLogger(Fleetd.class);
/** CB-637: how often the lead coordination loop looks for peer messages — see {@code Fleetd}'s own constant. */
private static final long LEAD_COORD_INTERVAL_MS = 3_000L;
private FleetdAssembly() {
}
static FleetdRuntime assembleAndStart(AssemblyInputs inputs, ResourcePorts ports) {
FleetConfig cfg = inputs.cfg();
ConfigRef config = inputs.config();
SubscriptionGuard guard = inputs.guard();
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
? Path.of(cfg.herdrSocket())
: UnixSocketHerdrClient.defaultSocketPath();
HerdrClient herdr = ports.connectHerdr(socket);
HerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()
? ports.connectHerdr(Path.of(cfg.memberHerdrSocket()))
: herdr;
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
target -> leadsRef.get().get().containsKey(target));
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
Map<String, FleetConfig.Profile> claudeProfiles = new LinkedHashMap<>();
Map<String, FleetConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
cfg.profiles().forEach((name, w) -> {
if (w.isOpenCode()) {
opencodeProfiles.put(name, w);
} else {
claudeProfiles.put(name, w);
}
});
List<HerdrPeerLauncher> adapters = new ArrayList<>();
// fleetd #175: the daemon's real ExhaustionSink can only be built once `sessions` exists
// (below), but `sessions` needs `workers`, which needs the adapters built right here — a
// genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now,
// pointed at the real one once it exists.
AtomicReference<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwardingExhaustionSink = Fleetd.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.
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
adapters.add(Fleetd.claudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
claudeProfiles, cfg, config));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(Fleetd.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
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
BackendQuarantine quarantine = BackendQuarantine.withEscalation(ports.nanoClock(),
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
// fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the
// launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in
// below (written on a classified backend error).
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(ports.nanoClock());
PeerLauncher workers = new CompositePeerLauncher(
adapters,
cfg.effectiveDefaultProfile(),
config,
profileName -> liveCountRef.get().apply(profileName),
quarantine,
outagePolicy);
// fleetd #422 follow-up: say which of the three model-gate states the daemon booted into.
log.info("model gate (fleetd #422): {}", Fleetd.modelGateCoverageLine(workers.modelGateState()));
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
// exists. Wait, then degrade rather than die: serving with /healthz reporting "degraded" is
// strictly more useful than exiting.
Fleetd.HerdrAwaitOutcome herdrOutcome = Fleetd.awaitHerdr(herdr, ports.nanoClock(), Fleetd::sleepHerdrPoll);
boolean herdrUp = Fleetd.logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
if (herdrUp) {
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
// with the previous process — reap those leaked orphans now, before we start serving.
workers.reapOrphanWorkers();
}
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
// CB-303 part 2: context cap is opt-in and disabled (0) when absent/null.
int contextCap = 0;
if (cfg.lifecycle() != null && cfg.lifecycle().contextCap() != null
&& cfg.lifecycle().contextCap() > 0) {
contextCap = cfg.lifecycle().contextCap();
}
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
SessionManager sessions = new SessionManager(workers,
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
ports.nanoClock(), contextCap, clearAfterTurn);
liveCountRef.set(profileName -> Fleetd.liveSessionCount(sessions.roster(), profileName));
// Idle-sleep guard: hold an OS-level assertion against idle sleep while at least one
// member is live. No-op (never constructed) off macOS or when idleSleepGuard.enabled is
// explicitly false; the mechanism itself is additionally a no-op if 'caffeinate' cannot be
// started, so this can never fail a spawn, a release, or startup.
boolean idleSleepGuardEnabled = cfg.idleSleepGuard() == null || cfg.idleSleepGuard().isEnabled();
final IdleSleepGuard idleSleepGuard;
if (idleSleepGuardEnabled) {
idleSleepGuard = new IdleSleepGuard(new CaffeinateSleepAssertionMechanism(), sessions::size);
sessions.onAcquire(_ -> idleSleepGuard.recheck());
sessions.onRelease(_ -> idleSleepGuard.recheck());
} else {
idleSleepGuard = null;
}
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled. FIRST of the
// recurring background loops to start — see this class's own javadoc for the full order.
final SessionReaper reaper;
if (cfg.lifecycle() != null
&& cfg.lifecycle().idleTtlSeconds() != null
&& cfg.lifecycle().idleTtlSeconds() > 0) {
reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds());
reaper.start();
} else {
reaper = null;
}
// CB-530: every pane the config names as a lead, merged from `leaders:` and the legacy
// singular pin. PrimaryRegistry below still tracks ONE terminal — it addresses the push
// loop's nudges, which need a single destination — so it keeps the legacy pin.
Map<String, String> leadTerminals = cfg.leaderTerminals();
if (leadTerminals.size() > 1) {
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
}
// CB-531/CB-579: discover leads by the tab labels the operator writes, one scanner per
// configured lead's own exact `tab:` label.
final Supplier<Map<String, String>> leads;
var leaders = cfg.fleet().leaders();
if (!leaders.isEmpty()) {
Map<String, String> tabToName = new LinkedHashMap<>();
leaders.forEach((name, leader) -> {
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
tabToName.put(leader.tab(), name);
}
});
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), ports.nanoClock());
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
tabToName.keySet(), scanIntervalSeconds);
} else {
leads = () -> leadTerminals;
}
leadsRef.set(leads);
// CB-558: start any declared lead that is not already running. After the scanner is built,
// and only when herdr answered — the launcher's whole safety property is that it can count
// live leads first, and must never guess and risk a second orchestrator.
if (herdrUp && !leaders.isEmpty()) {
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
if (launched > 0) {
log.info("lead auto-launch: {} lead(s) started", launched);
}
}
// CB-548: config-declared architect slots. Nothing here spawns a slot; the terminal → slot
// binding is owned by the registry and empty at startup.
MemberRegistry members = MemberRegistry.live(() -> config.get().fleet());
sessions.setMemberLifecycle(members);
if (!members.slots().isEmpty()) {
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
+ "spawn lifecycle binds a live terminal to it)",
members.slots().size(), members.slots().keySet());
}
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
Rendezvous rendezvous = new Rendezvous();
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
LiveExhaustedPatterns liveExhaustedPatterns = Fleetd.liveExhaustedPatterns(config);
ExhaustedPatternLookup exhaustedPatterns = Fleetd.exhaustedPatternLookup(sessions::roster, liveExhaustedPatterns);
// The startup coverage line still reports the boot-time snapshot only.
Set<String> exhaustedConfiguredAtStartup = cfg.profiles().entrySet().stream()
.filter(e -> e.getValue().hasExhaustedPattern())
.map(Map.Entry::getKey)
.collect(Collectors.toCollection(LinkedHashSet::new));
log.info("backend-exhausted classification (CB-578 stage A): {}",
Fleetd.exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedConfiguredAtStartup));
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
// rather than handing it back as a real answer.
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasErrorPattern()) {
errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern()));
}
});
BackendErrorPatternLookup backendErrorPatterns = Fleetd.backendErrorPatternLookup(sessions::roster,
errorPatternsByProfile);
log.info("backend-error classification (fleetd #201 Unit 5): {}",
Fleetd.errorPatternCoverageLine(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
// CREDENTIAL, so a profile sharing that credential is refused too, not just the one that
// happened to report it.
Map<String, String> quarantineReasonByCredential = new ConcurrentHashMap<>();
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
// now that `sessions` exists to resolve target -> session -> profile.
ExhaustionSink exhaustionSink = Fleetd.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
// holder set once `pushLoop` exists, read lazily from inside the lambda built here.
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>();
BackendErrorSink backendErrorSink = Fleetd.backendErrorSink(sessions, () -> config.get().profiles(),
outagePolicy, pushLoopRef::get);
AgentControl agents = router.memberAgents();
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
exhaustionSink, backendErrorPatterns, backendErrorSink, ports.nanoClock(),
Fleetd.worktreeBranchLookup(sessions::roster));
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = Fleetd.turnListener(completion, sessions);
Predicate<String> deliverable = Fleetd.deliverableTo(presence, leads);
// 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.
Injector injector = new Injector(router, turnListener, deliverable,
presence::forget, Fleetd.turnRegistrar(completion));
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start(); // SECOND of the recurring background loops to start, after the reaper.
// 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 = Fleetd.selectReplyInbox(cfg.broker(), ports.environment(),
ports.replyInboxOpener());
// CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate
// broker from the reply inbox by design. Absent a coordinator: block this is null and every
// lead path below is simply not wired, exactly the behaviour before this ticket. It owns a
// broker connection, so keep the reference for the ordered shutdown hook.
final LeadMailbox leadMailbox = Fleetd.openLeadMailbox(cfg.coordinator(), ports.environment(),
ports.leadMailboxOpener());
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs.
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
+ "from leaders:/leadScan:, and reply nudges follow the lead that delegated. "
+ "It still works, and is still the fallback nudge destination when a restart "
+ "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid.");
}
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
var pushScheduler = ports.newScheduler("bridge-push-");
// CB-502: the registry is built before the service and the push loop so send/reply outcomes
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
// above at the real push loop, now that it exists.
pushLoopRef.set(pushLoop);
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
// THIRD of the recurring background loops to start (optional).
final LeadHeartbeatLoop heartbeat;
var heartbeatScheduler = ports.newScheduler("bridge-heartbeat-");
if (cfg.leadHeartbeat() != null) {
var hb = cfg.leadHeartbeat();
var leadContextGauge = new LeadContextGauge();
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
pushLoop, heartbeatScheduler, ports.nanoClock(),
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
metrics,
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders)),
Boolean.TRUE.equals(hb.contextHighNudge()));
heartbeat.start();
} else {
heartbeat = null;
heartbeatScheduler.shutdownNow();
}
// fleetd #480: lead rollover. Opt-in; absent `leadRollover:` this is never constructed.
LeadRollover leadRollover = Fleetd.leadRollover(cfg, router.leadAgents(), config, leads);
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
pushLoop, metrics);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
// FOURTH of the recurring background loops to start (optional).
final FleetHealthMonitor healthMonitor;
var healthScheduler = ports.newScheduler("bridge-health-");
if (cfg.health() != null && cfg.health().isEnabled()) {
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
ports.nanoClock(), ports.wallClockNanos(),
cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), Fleetd.healthFailTarget(messages));
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
log.warn("fleet health: {} (no notification sink configured)", coverage);
} else {
log.info("fleet health: {}", coverage);
}
healthMonitor.start();
} else {
healthMonitor = null;
healthScheduler.shutdownNow();
}
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
sessions.onAcquire(replyInbox::own);
// CB-516: releasing a worker must fail whatever send was waiting on it.
sessions.onRelease(Fleetd.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.
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
final CallerResolver callers;
if (cfg.auth().tokenMode()) {
String token = ports.environment().get(cfg.auth().tokenEnv());
if (token == null || token.isBlank()) {
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
+ " is unset or empty — export it before starting fleetd");
}
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
log.info("auth: token mode (bearer required for non-worker callers, env {})",
cfg.auth().tokenEnv());
} else {
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
FleetMcp.QuarantineSource quarantineSource = Fleetd.quarantineSource(config, quarantine,
liveExhaustedPatterns, quarantineReasonByCredential);
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp.LoopHealthSource loopHealth = Fleetd.loopHealthSource(poller, reaper);
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
Fleetd.capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
Fleetd.healthCoverageSource(config),
loopHealth,
quarantineSource,
leadMailbox,
outageSource,
new FleetMcp.LeadSeatSource(Fleetd.leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
Fleetd.leadConfigDirSource(() -> config.get().profiles(), leaders),
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers(),
leadRollover);
// CB-637: the receive half. Only constructed when a lead mailbox actually opened.
final LeadCoordLoop leadCoordLoop;
final ScheduledExecutorService leadCoordSchedulerRef;
if (leadMailbox != null) {
var leadCoordScheduler = ports.newScheduler("bridge-leadcoord-");
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
LEAD_COORD_INTERVAL_MS);
leadCoordLoop.start(); // FIFTH of the recurring background loops to start (optional).
leadCoordSchedulerRef = leadCoordScheduler;
} else {
leadCoordLoop = null;
leadCoordSchedulerRef = null;
}
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed.
// SIXTH of the recurring background loops to start (optional).
final ConfigWatcher configWatcher;
if (cfg.configReload() != null && cfg.configReload().isEnabled()) {
configWatcher = new ConfigWatcher(config, cfg.configReload().intervalSeconds());
configWatcher.start();
} else {
configWatcher = null;
}
// CB-303 part 3: single ordered shutdown hook — see FleetdRuntime#close() for the statements
// this used to be. Registered here, at the exact point `main` used to register it: after
// configWatcher, before the Javalin app exists (see this class's own javadoc for why).
FleetdRuntime runtime = new FleetdRuntime(cfg, sessions, router, poller, messages, pushLoop, heartbeat,
leadCoordLoop, leadCoordSchedulerRef, healthMonitor, configWatcher, mcp, reaper, idleSleepGuard,
replyInbox, leadMailbox, completion, injector);
ports.addShutdownHook(runtime::close);
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable,
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
quarantineSource, outageSource, loopHealth).build();
runtime.attachApp(app);
ports.startHttp(app, cfg.bind().host(), cfg.bind().port()); // HTTP starts LAST, always.
log.info("fleetd listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
return runtime;
}
}
@@ -0,0 +1,166 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigWatcher;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.health.FleetHealthMonitor;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.msg.LeadCoordLoop;
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.power.IdleSleepGuard;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.SessionReaper;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.concurrent.ScheduledExecutorService;
/**
* fleetd #612 Unit A: the assembled daemon. {@link FleetdAssembly#assembleAndStart} builds exactly
* one of these and hands it to {@code ResourcePorts.addShutdownHook}; {@link #close()} is the
* single ordered shutdown, moved verbatim out of {@code Fleetd.main}'s old shutdown-hook
* {@code Thread} body — same statements, same order, see that method's javadoc.
*
* <p><strong>Owns the real production objects, never a copy.</strong> Every package-private
* accessor below returns the identical instance the running daemon is using. That is the entire
* point of this class existing (see fleetd #612's problem statement): a test that inspected a
* snapshot built alongside the real objects could pass while production silently received
* something else — the exact shape of the #602/#606 defect this ticket exists to stop from
* recurring one call site at a time. Nothing here is rebuilt or copied for a test's benefit.
*/
final class FleetdRuntime implements AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(FleetdRuntime.class);
private final FleetConfig cfg;
private final SessionManager sessions;
private final HerdrRouter router;
private final StatusPoller poller;
private final MessageService messages;
private final ReplyPushLoop pushLoop;
private final LeadHeartbeatLoop heartbeat; // nullable — leadHeartbeat: opt-in
private final LeadCoordLoop leadCoordLoop; // nullable — coordinator: opt-in
private final ScheduledExecutorService leadCoordScheduler; // nullable, paired with leadCoordLoop
private final FleetHealthMonitor healthMonitor; // nullable — health.enabled opt-in
private final ConfigWatcher configWatcher; // nullable — configReload.enabled opt-in
private final FleetMcp mcp;
private final SessionReaper reaper; // nullable — lifecycle.idleTtlSeconds opt-in
private final IdleSleepGuard idleSleepGuard; // nullable — idleSleepGuard.enabled: false
private final ReplyInbox replyInbox;
private final LeadMailbox leadMailbox; // nullable — coordinator: opt-in
private final CompletionResolver completion;
private final Injector injector;
/**
* Not final: {@code Fleetd.main}'s shutdown hook was registered <em>before</em> the Javalin
* {@code FleetApp} was built and started — see {@link FleetdAssembly#assembleAndStart}'s javadoc
* for why that order could not be preserved AND have this constructor take {@code app}. {@link
* #attachApp} is called immediately after the real app is built, still before HTTP starts
* listening, so this is set long before any test or caller could observe it unset.
*/
private Javalin app;
FleetdRuntime(FleetConfig cfg, SessionManager sessions, HerdrRouter router, StatusPoller poller,
MessageService messages, ReplyPushLoop pushLoop, LeadHeartbeatLoop heartbeat,
LeadCoordLoop leadCoordLoop, ScheduledExecutorService leadCoordScheduler,
FleetHealthMonitor healthMonitor, ConfigWatcher configWatcher, FleetMcp mcp,
SessionReaper reaper, IdleSleepGuard idleSleepGuard, ReplyInbox replyInbox,
LeadMailbox leadMailbox, CompletionResolver completion, Injector injector) {
this.cfg = cfg;
this.sessions = sessions;
this.router = router;
this.poller = poller;
this.messages = messages;
this.pushLoop = pushLoop;
this.heartbeat = heartbeat;
this.leadCoordLoop = leadCoordLoop;
this.leadCoordScheduler = leadCoordScheduler;
this.healthMonitor = healthMonitor;
this.configWatcher = configWatcher;
this.mcp = mcp;
this.reaper = reaper;
this.idleSleepGuard = idleSleepGuard;
this.replyInbox = replyInbox;
this.leadMailbox = leadMailbox;
this.completion = completion;
this.injector = injector;
}
/** See the {@link #app} field doc for why this is a late-bound setter rather than a constructor arg. */
void attachApp(Javalin app) {
this.app = app;
}
// --- package-private accessors: the SAME instances this runtime owns, never a copy ---------
SessionManager sessions() { return sessions; }
HerdrRouter router() { return router; }
StatusPoller poller() { return poller; }
MessageService messages() { return messages; }
ReplyPushLoop pushLoop() { return pushLoop; }
LeadHeartbeatLoop heartbeat() { return heartbeat; }
LeadCoordLoop leadCoordLoop() { return leadCoordLoop; }
FleetHealthMonitor healthMonitor() { return healthMonitor; }
ConfigWatcher configWatcher() { return configWatcher; }
FleetMcp mcp() { return mcp; }
SessionReaper reaper() { return reaper; }
IdleSleepGuard idleSleepGuard() { return idleSleepGuard; }
ReplyInbox replyInbox() { return replyInbox; }
LeadMailbox leadMailbox() { return leadMailbox; }
CompletionResolver completion() { return completion; }
Injector injector() { return injector; }
Javalin app() { return app; }
/**
* CB-303 part 3: the single ordered shutdown. Moved verbatim out of {@code Fleetd.main}'s
* shutdown-hook {@code Thread} body (fleetd #612 Unit A) — drain sessions first while herdr is
* still open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close
* herdr last, exactly as before. {@code Fleetd.main} never calls this directly; it hands the
* reference to {@code ResourcePorts.addShutdownHook} the moment this runtime exists, the same
* point it used to register the hook {@code Thread} itself.
*/
@Override
public void close() {
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
poller.stop();
messages.close();
pushLoop.close();
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
if (leadCoordLoop != null) leadCoordLoop.close(); // CB-637: stop delivering peer-lead messages
if (leadCoordScheduler != null) leadCoordScheduler.shutdownNow();
if (healthMonitor != null) healthMonitor.stop();
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
mcp.close();
if (reaper != null) reaper.stop();
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
// drained every session (and each release already drove the live count to 0, which
// releases the guard's assertion on its own) — this is the backstop for a drain that was
// itself interrupted or threw, so no caffeinate child ever outlives the daemon.
if (idleSleepGuard != null) idleSleepGuard.close();
// Release the broker connection last among message resources (no-op for the in-memory inbox).
if (replyInbox instanceof AutoCloseable closeable) {
try {
closeable.close();
} catch (Exception e) {
log.debug("reply inbox close: {}", e.toString());
}
}
// CB-637: the coordination connection goes with it — after the loop that reads it has
// stopped, so no tick can be mid-ack against a closed channel.
if (leadMailbox != null) {
try {
leadMailbox.close();
} catch (Exception e) {
log.debug("lead mailbox close: {}", e.toString());
}
}
router.close();
}
}
@@ -0,0 +1,67 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.HerdrClient;
import io.javalin.Javalin;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
/**
* fleetd #612 Unit A: every boot-time side effect {@link FleetdAssembly#assembleAndStart} performs
* that a real daemon must do for real, and a test must not — read the process environment, connect
* a herdr client, open a broker (the reply inbox, the lead mailbox), read a clock, start a
* background scheduler, register the JVM shutdown hook, and bind the HTTP server.
*
* <p>{@link #system()} is the one production implementation ({@code SystemResourcePorts}), wired
* verbatim from what {@code Fleetd.main} used to call directly at each of these call sites. A test
* builds its own implementation instead of receiving an inert default from this interface —
* deliberately, there is no {@code ResourcePorts.none()}. fleetd #612's whole problem is a call
* site quietly swapped for an inert variant that still compiles; adding one here, even for tests,
* would hand a future edit to {@code FleetdAssembly} the exact compiling substitute this ticket
* exists to rule out. A test that wants an inert resource writes its own fake and owns that
* decision explicitly.
*/
public interface ResourcePorts {
/** The process environment. Production: {@link System#getenv()}. */
Map<String, String> environment();
/**
* Connect a herdr client bound to {@code socketPath}. Production returns a real
* {@code UnixSocketHerdrClient} — connection-per-call, so this itself never touches the socket.
*/
HerdrClient connectHerdr(Path socketPath);
/** The reply-inbox AMQP opener (CB-307). Production: {@link Fleetd#replyInboxOpener()}. */
Fleetd.AmqpOpener replyInboxOpener();
/** The lead-mailbox AMQP opener (CB-637). Production: {@link Fleetd#leadMailboxOpener()}. */
Fleetd.LeadMailboxOpener leadMailboxOpener();
/** A monotonic elapsed-time clock. Production: {@link System#nanoTime()}. */
LongSupplier nanoClock();
/**
* A wall-clock reading, in nanoseconds. Production: {@code System.currentTimeMillis()}
* converted to nanoseconds. Kept separate from {@link #nanoClock()} because {@link
* dev.ltms.fleet.health.FleetHealthMonitor} needs both — one monotonic clock for elapsed-time
* decisions, one wall clock to detect and correct for a macOS sleep freezing the monotonic one.
*/
LongSupplier wallClockNanos();
/** A dedicated single-thread scheduler; production names its (virtual) thread {@code purpose}. */
ScheduledExecutorService newScheduler(String purpose);
/** Register a JVM shutdown hook that runs {@code hook} on JVM exit. */
void addShutdownHook(Runnable hook);
/** Bind and start the HTTP server. */
void startHttp(Javalin app, String host, int port);
/** The real production ports: a live herdr socket, a live broker, real threads, a real bind. */
static ResourcePorts system() {
return new SystemResourcePorts();
}
}
@@ -0,0 +1,66 @@
package dev.ltms.fleet;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
import io.javalin.Javalin;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
/**
* fleetd #612 Unit A: the one production {@link ResourcePorts} — every method here is the exact
* call {@code Fleetd.main} used to make directly at each of these sites before this ticket.
* Package-private: obtained only through {@link ResourcePorts#system()}.
*/
final class SystemResourcePorts implements ResourcePorts {
@Override
public Map<String, String> environment() {
return System.getenv();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return UnixSocketHerdrClient.connect(socketPath, new ObjectMapper());
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return Fleetd.replyInboxOpener();
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return Fleetd.leadMailboxOpener();
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor(r -> Thread.ofVirtual().name(purpose).unstarted(r));
}
@Override
public void addShutdownHook(Runnable hook) {
Runtime.getRuntime().addShutdownHook(new Thread(hook));
}
@Override
public void startHttp(Javalin app, String host, int port) {
app.start(host, port);
}
}
@@ -0,0 +1,258 @@
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.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
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.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 Unit A: proves {@link FleetdAssembly#assembleAndStart} — not a copy of its logic —
* against a real {@link FleetConfig}, a {@link FakeHerdr} and a fake {@link ResourcePorts}, with no
* real herdr socket, no real broker, and no real HTTP bind.
*
* <p><strong>The canonical order this test asserts against was recorded from {@code Fleetd.java}
* BEFORE any code moved</strong> (fleetd #612 Unit A's mandated order of work), by reading the
* original {@code main}'s body and its shutdown-hook {@code Thread}:
*
* <p>Start order: {@code SessionReaper.start()} → {@code StatusPoller.start()} →
* {@code LeadHeartbeatLoop.start()} (opt-in) → {@code FleetHealthMonitor.start()} (opt-in) →
* {@code LeadCoordLoop.start()} (opt-in) → {@code ConfigWatcher.start()} (opt-in) →
* {@code app.start()} (HTTP), always last.
*
* <p>Close order (from the original shutdown hook body): {@code sessions.close(drainTimeoutSeconds)}
* → {@code poller.stop()} → {@code messages.close()} → {@code pushLoop.close()} →
* {@code heartbeat.close()} (if present) → {@code leadCoordLoop.close()} (if present) →
* {@code leadCoordScheduler.shutdownNow()} (if present) → {@code healthMonitor.stop()} (if present)
* → {@code configWatcher.stop()} (if present) → {@code mcp.close()} → {@code reaper.stop()} (if
* present) → {@code idleSleepGuard.close()} (if present) → {@code replyInbox.close()} (if
* {@code AutoCloseable}) → {@code leadMailbox.close()} (if present) → {@code router.close()}.
*
* <p>This test's config deliberately leaves {@code coordinator:} unset, so {@code leadMailbox} and
* {@code leadCoordLoop} stay {@code null} throughout — the lead-mailbox resource-ledger criterion is
* NOT exercised here; see the class-level caveat in the implementer's hand-off. {@code
* idleSleepGuard.enabled: false} is set for the same kind of reason: it would otherwise try to spawn
* a real {@code caffeinate} subprocess, which is not one of the resources the ticket's acceptance
* criteria names (scheduler/inbox/mailbox/loop/MCP server/herdr router).
*/
class FleetdAssemblyLifecycleTest {
/** The one {@link HerdrClient} both {@code herdrSocket} and {@code memberHerdrSocket} resolve to. */
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
final List<String> ledger = new CopyOnWriteArrayList<>();
final List<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
final List<String> schedulerPurposes = new ArrayList<>();
Runnable shutdownHook;
Javalin startedApp;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
ledger.add("connectHerdr");
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
ledger.add("replyInboxOpener");
return (uri, prefetch) -> replyInbox;
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
// Never invoked: this test's config has no `coordinator:` block, so
// Fleetd.openLeadMailbox returns null before calling the opener at all.
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
ledger.add("newScheduler:" + purpose);
schedulerPurposes.add(purpose);
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
schedulers.add(scheduler);
return scheduler;
}
@Override
public void addShutdownHook(Runnable hook) {
ledger.add("addShutdownHook");
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Deliberately never call app.start(host, port): no real HTTP bind in a unit test.
ledger.add("startHttp");
this.startedApp = app;
}
}
/** A fake {@link ReplyInbox} that is also {@link AutoCloseable}, so the ledger can prove it closes. */
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
volatile boolean closed = false;
@Override
public void own(String target) {
}
@Override
public void release(String target) {
}
@Override
public void publish(String target, String msgId, String content) {
}
@Override
public List<InboxMessage> peek(String target) {
return List.of();
}
@Override
public boolean ack(String target, String msgId) {
return false;
}
@Override
public void close() {
closed = true;
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
lifecycle:
idleTtlSeconds: 600
health:
enabled: true
intervalSeconds: 30
leadHeartbeat:
idleAfterSeconds: 600
backoffMs: 15000
quietNudgeCap: 5
configReload:
enabled: true
intervalSeconds: 30
idleSleepGuard:
enabled: false
broker:
uri: "amqp://fake-test-broker/vh"
""");
return FleetConfig.load(f);
}
@Test
void assemblesTheRealBootGraphWithoutTouchingAnyRealSocketBrokerOrPort(@TempDir Path dir)
throws Exception {
FleetConfig cfg = writeConfig(dir);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
// --- start order: the herdr connect, then the recurring background loops, in the recorded
// order, then the shutdown hook is registered, then (finally) HTTP "starts" -------------
assertTrue(ports.ledger.indexOf("connectHerdr") < ports.ledger.indexOf("newScheduler:bridge-push-"),
"herdr must connect before the push scheduler is created: " + ports.ledger);
assertEquals(List.of("bridge-push-", "bridge-heartbeat-", "bridge-health-"), ports.schedulerPurposes,
"the three always-created schedulers must be requested in exactly this order: "
+ ports.schedulerPurposes);
assertTrue(ports.ledger.indexOf("addShutdownHook") < ports.ledger.indexOf("startHttp"),
"the shutdown hook must be registered before HTTP starts — the one statement that "
+ "could not be reordered without changing FleetdRuntime's constructor shape, "
+ "see FleetdAssembly's javadoc: " + ports.ledger);
assertEquals(ports.ledger.size() - 1, ports.ledger.indexOf("startHttp"),
"HTTP must start LAST of everything this fake observes: " + ports.ledger);
assertNotNull(ports.startedApp, "FleetdAssembly must have built and handed off a real FleetApp");
// No real HTTP bind and no real herdr socket: this call returning at all, plus the ledger
// above, is the proof — a real bind or a real UnixSocketHerdrClient.connect would have
// thrown or hung against the sockets/ports this test never opened.
assertNotNull(runtime.app(), "FleetdRuntime must own the same Javalin app that was built");
assertEquals(ports.startedApp, runtime.app(), "attachApp must hand FleetdRuntime the SAME instance");
// --- the runtime owns the real, live objects — not a copy --------------------------------
assertEquals(LoopWatchdog.State.RUNNING, runtime.reaper().health(),
"SessionReaper must be running: lifecycle.idleTtlSeconds is configured");
assertEquals(LoopWatchdog.State.RUNNING, runtime.poller().health(), "StatusPoller must be running");
assertNotNull(runtime.heartbeat(), "leadHeartbeat: is configured, so the loop must be built and started");
assertNotNull(runtime.healthMonitor(), "health.enabled: true, so the monitor must be built and started");
assertNotNull(runtime.configWatcher(), "configReload.enabled: true, so the watcher must be built and started");
// Documented gap (see class javadoc): no coordinator: block, so these stay null.
assertEquals(null, runtime.leadCoordLoop(), "no coordinator: block — leadCoordLoop must stay unbuilt");
assertEquals(null, runtime.leadMailbox(), "no coordinator: block — leadMailbox must stay unbuilt");
assertEquals(runtime.replyInbox(), ports.replyInbox,
"FleetdRuntime must own the exact ReplyInbox instance this fake's AmqpOpener returned");
assertFalse(ports.herdr.closed, "herdr must still be open while the daemon is running");
assertFalse(ports.replyInbox.closed, "the reply inbox must still be open while the daemon is running");
for (ScheduledExecutorService scheduler : ports.schedulers) {
assertFalse(scheduler.isShutdown(), "a scheduler must still be running while the daemon is up");
}
// --- close order: invoke the captured shutdown-hook Runnable directly (no real JVM shutdown
// happens in a unit test) and prove every resource this fake can observe is released --------
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
ports.shutdownHook.run();
assertEquals(LoopWatchdog.State.STOPPED, runtime.reaper().health(), "SessionReaper must stop on close");
assertEquals(LoopWatchdog.State.STOPPED, runtime.poller().health(), "StatusPoller must stop on close");
assertTrue(ports.herdr.closed, "router.close() must close the herdr client last");
assertTrue(ports.replyInbox.closed, "the AutoCloseable reply inbox must be closed");
for (int i = 0; i < ports.schedulers.size(); i++) {
assertTrue(ports.schedulers.get(i).isShutdown(),
"scheduler for '" + ports.schedulerPurposes.get(i) + "' must be shut down by close(): "
+ "LeadHeartbeatLoop/FleetHealthMonitor/ReplyPushLoop each call "
+ "scheduler.shutdownNow() on the exact instance ports.newScheduler(...) handed them");
}
// Calling the captured hook a second time must never happen for a real JVM shutdown hook,
// but nothing above should have thrown — that already proves every accessed field's close()
// tolerated running once, in the recorded order, without an exception escaping.
}
}
@@ -456,7 +456,11 @@ public final class FakeHerdr implements HerdrClient {
}
}
/** fleetd #612 Unit A: set by {@link #close()} so a test can prove a router/client actually closed this. */
public volatile boolean closed = false;
@Override
public void close() {
closed = true;
}
}