diff --git a/fleetd/src/main/java/dev/ltms/fleet/AssemblyInputs.java b/fleetd/src/main/java/dev/ltms/fleet/AssemblyInputs.java new file mode 100644 index 0000000..1b78ab1 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/AssemblyInputs.java @@ -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) { +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index f9b42d6..43896bf 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -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>> 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 claudeProfiles = new LinkedHashMap<>(); - Map opencodeProfiles = new LinkedHashMap<>(); - cfg.profiles().forEach((name, w) -> { - if (w.isOpenCode()) { - opencodeProfiles.put(name, w); - } else { - claudeProfiles.put(name, w); - } - }); - List 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 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> 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 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> leads; - var leaders = cfg.fleet().leaders(); - if (!leaders.isEmpty()) { - Map 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 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 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 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 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 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. + * + *

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) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java new file mode 100644 index 0000000..ebba783 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java @@ -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 the same statements {@code main} used to run + * inline, 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. + * + *

Construction and start order is preserved exactly, on purpose. 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. + * + *

One statement could not move without reordering startup. {@code Fleetd.main} + * registered its shutdown hook before 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. + * + *

No inert variant. 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>> 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 claudeProfiles = new LinkedHashMap<>(); + Map opencodeProfiles = new LinkedHashMap<>(); + cfg.profiles().forEach((name, w) -> { + if (w.isOpenCode()) { + opencodeProfiles.put(name, w); + } else { + claudeProfiles.put(name, w); + } + }); + List 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 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> 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 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> leads; + var leaders = cfg.fleet().leaders(); + if (!leaders.isEmpty()) { + Map 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 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 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 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 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 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; + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/FleetdRuntime.java b/fleetd/src/main/java/dev/ltms/fleet/FleetdRuntime.java new file mode 100644 index 0000000..4ae5e44 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/FleetdRuntime.java @@ -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. + * + *

Owns the real production objects, never a copy. 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 before 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(); + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/ResourcePorts.java b/fleetd/src/main/java/dev/ltms/fleet/ResourcePorts.java new file mode 100644 index 0000000..bec9d42 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/ResourcePorts.java @@ -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. + * + *

{@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 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(); + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/SystemResourcePorts.java b/fleetd/src/main/java/dev/ltms/fleet/SystemResourcePorts.java new file mode 100644 index 0000000..7ce5ddd --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/SystemResourcePorts.java @@ -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 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); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLifecycleTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLifecycleTest.java new file mode 100644 index 0000000..8f8f488 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdAssemblyLifecycleTest.java @@ -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. + * + *

The canonical order this test asserts against was recorded from {@code Fleetd.java} + * BEFORE any code moved (fleetd #612 Unit A's mandated order of work), by reading the + * original {@code main}'s body and its shutdown-hook {@code Thread}: + * + *

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. + * + *

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()}. + * + *

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 ledger = new CopyOnWriteArrayList<>(); + final List schedulers = new CopyOnWriteArrayList<>(); + final List schedulerPurposes = new ArrayList<>(); + Runnable shutdownHook; + Javalin startedApp; + + @Override + public Map 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 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. + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/herdr/FakeHerdr.java b/fleetd/src/test/java/dev/ltms/fleet/herdr/FakeHerdr.java index 5d3fbc4..a123aa4 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/herdr/FakeHerdr.java +++ b/fleetd/src/test/java/dev/ltms/fleet/herdr/FakeHerdr.java @@ -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; } }