From 6af87b6ad616e14a0fff0cd2753e3abd9fec9a4c Mon Sep 17 00:00:00 2001 From: Ha Trong Dai Date: Fri, 28 Aug 2026 09:42:45 +0700 Subject: [PATCH] CB-185: share routed herdr controls --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 20 ++++++++++--------- .../dev/ltms/fleet/inject/StatusPoller.java | 4 ++-- 2 files changed, 13 insertions(+), 11 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 8ac43a6..958866b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -153,6 +153,9 @@ public final class Fleetd { 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. @@ -170,14 +173,14 @@ public final class Fleetd { // 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(new ClaudeCodeLauncher(new AgentControl(memberHerdr), new WorkspaceControl(memberHerdr), guard, + adapters.add(new ClaudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard, claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(), () -> config.get().fleet(), () -> config.get().memberCredentials())); } if (!opencodeProfiles.isEmpty()) { - adapters.add(new OpenCodeLauncher(new AgentControl(memberHerdr), new WorkspaceControl(memberHerdr), + adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(), opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(), () -> config.get().fleet(), @@ -280,13 +283,14 @@ public final class Fleetd { } 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(new AgentControl(herdr), new WorkspaceControl(herdr), cfg).ensureLeads(); + int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads(); if (launched > 0) { log.info("lead auto-launch: {} lead(s) started", launched); } @@ -342,8 +346,6 @@ public final class Fleetd { + "BACKEND_EXHAUSTED): {}", credentialId, cfg.quarantineCooldownSeconds(), profile.profile(), reason); }); - HerdrRouter router = new HerdrRouter(herdr, memberHerdr, - target -> leads.get().containsKey(target)); AgentControl agents = router.memberAgents(); CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink); // CB-113: deliver only to an available worker (its MCP is connected), never its boot window. @@ -425,7 +427,7 @@ public final class Fleetd { // 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, agents, replyInbox, + var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox, pushScheduler, maxReminders, backoffMs, metrics); // 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. @@ -435,7 +437,7 @@ public final class Fleetd { Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r)); if (cfg.leadHeartbeat() != null) { var hb = cfg.leadHeartbeat(); - heartbeat = new LeadHeartbeatLoop(primaryRegistry, agents, replyInbox, sessions::roster, + heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster, pushLoop, heartbeatScheduler, System::nanoTime, TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(), metrics); @@ -497,7 +499,7 @@ public final class Fleetd { // 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), new LsofPeerPidLookup(), new LsofProcessCwdLookup()); + new PaneLocator(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. @@ -543,7 +545,7 @@ public final class Fleetd { if (leadMailbox != null) { var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r -> Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r)); - leadCoordLoop = new LeadCoordLoop(leadMailbox, agents, leads, leadCoordScheduler, + leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler, LEAD_COORD_INTERVAL_MS); leadCoordLoop.start(); leadCoordSchedulerRef = leadCoordScheduler; diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java index 87180ea..5dda76c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java @@ -47,7 +47,7 @@ public final class StatusPoller { this.agents = null; this.router = router; this.injector = injector; - this.refiner = null; + this.refiner = new StatusRefiner(router.memberAgents()); this.intervalMillis = intervalMillis; } @@ -68,7 +68,7 @@ public final class StatusPoller { // herdr's agent_status can misreport a settled worker as `unknown`; refine it // against the pane content before it drives delivery/completion (CB-115). AgentControl control = router != null ? router.agentsFor(target) : agents; - AgentStatus status = router != null ? control.status(target) : refiner.refine(target, control.status(target)); + AgentStatus status = refiner.refine(target, control.status(target)); injector.onStatus(target, status); } catch (HerdrException e) { // The worker's agent is gone — stop trying and unblock its waiters.