CB-185: share routed herdr controls
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m40s

This commit is contained in:
Ha Trong Dai
2026-08-28 09:42:45 +07:00
parent fc655e78c2
commit 6af87b6ad6
2 changed files with 13 additions and 11 deletions
@@ -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<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
target -> leadsRef.get().get().containsKey(target));
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
@@ -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;
@@ -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.