From fc655e78c29db85fd849bd3e69a9e5c75e7d609f Mon Sep 17 00:00:00 2001 From: Ha Trong Dai Date: Fri, 28 Aug 2026 09:37:44 +0700 Subject: [PATCH] CB-185: route members to separate herdr --- fleetd/fleetd.example.yaml | 3 ++ .../src/main/java/dev/ltms/fleet/Fleetd.java | 23 ++++++----- .../java/dev/ltms/fleet/config/ConfigRef.java | 5 ++- .../dev/ltms/fleet/config/FleetConfig.java | 26 ++++++------ .../dev/ltms/fleet/herdr/HerdrRouter.java | 40 +++++++++++++++++++ .../java/dev/ltms/fleet/inject/Injector.java | 20 +++++++++- .../dev/ltms/fleet/inject/StatusPoller.java | 14 ++++++- .../dev/ltms/fleet/herdr/HerdrRouterTest.java | 31 ++++++++++++++ 8 files changed, 137 insertions(+), 25 deletions(-) create mode 100644 fleetd/src/main/java/dev/ltms/fleet/herdr/HerdrRouter.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/herdr/HerdrRouterTest.java diff --git a/fleetd/fleetd.example.yaml b/fleetd/fleetd.example.yaml index f878fed..21c718e 100644 --- a/fleetd/fleetd.example.yaml +++ b/fleetd/fleetd.example.yaml @@ -114,6 +114,9 @@ bind: # (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}). herdrSocket: ~/.config/herdr/herdr.sock +# Optional socket for member panes. Omit this to use herdrSocket for both leads and members. +# memberHerdrSocket: /Users/member/.config/herdr/herdr.sock + # How member sessions are spawned. Define one or more named profiles (backends) under # `profiles`; each key is the profile name (also the ccs profile). A profile says only WHICH # BACKEND — model, CLI adapter, credentials, cost. It says nothing about what a member spawned on diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 1c265eb..8ac43a6 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -7,6 +7,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.HerdrClient; import dev.ltms.fleet.herdr.HerdrException; +import dev.ltms.fleet.herdr.HerdrRouter; import dev.ltms.fleet.herdr.LeadTabScanner; import dev.ltms.fleet.lead.LeadLauncher; import dev.ltms.fleet.herdr.PaneLocator; @@ -149,9 +150,9 @@ public final class Fleetd { : UnixSocketHerdrClient.defaultSocketPath(); UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper()); - - AgentControl agents = new AgentControl(herdr); - WorkspaceControl spaces = new WorkspaceControl(herdr); + UnixSocketHerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank() + ? UnixSocketHerdrClient.connect(Path.of(cfg.memberHerdrSocket()), new com.fasterxml.jackson.databind.ObjectMapper()) + : herdr; // 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. @@ -169,14 +170,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(agents, spaces, guard, + adapters.add(new ClaudeCodeLauncher(new AgentControl(memberHerdr), new WorkspaceControl(memberHerdr), guard, claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(), () -> config.get().fleet(), () -> config.get().memberCredentials())); } if (!opencodeProfiles.isEmpty()) { - adapters.add(new OpenCodeLauncher(agents, spaces, + adapters.add(new OpenCodeLauncher(new AgentControl(memberHerdr), new WorkspaceControl(memberHerdr), opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(), () -> config.get().fleet(), @@ -271,6 +272,7 @@ public final class Fleetd { // 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)", @@ -284,7 +286,7 @@ public final class Fleetd { // 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(agents, spaces, cfg).ensureLeads(); + int launched = new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg).ensureLeads(); if (launched > 0) { log.info("lead auto-launch: {} lead(s) started", launched); } @@ -340,6 +342,9 @@ 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. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY. @@ -381,9 +386,9 @@ public final class Fleetd { } }; Predicate deliverable = deliverableTo(presence, leads); - Injector injector = new Injector(agents, turnListener, deliverable, + Injector injector = new Injector(router, turnListener, deliverable, presence::forget); - StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS); + 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 @@ -589,7 +594,7 @@ public final class Fleetd { log.debug("lead mailbox close: {}", e.toString()); } } - herdr.close(); + router.close(); })); Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(), diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java b/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java index 45d6360..ce13fc1 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/ConfigRef.java @@ -71,7 +71,7 @@ public final class ConfigRef implements Supplier { /** Keys that cannot change under a running daemon — see the class doc. */ private static final Set COLD_KEYS = - Set.of("bind", "herdrSocket", "broker", "auth"); + Set.of("bind", "herdrSocket", "memberHerdrSocket", "broker", "auth"); private final Path path; private final AtomicReference current; @@ -190,6 +190,9 @@ public final class ConfigRef implements Supplier { if (!Objects.equals(old.herdrSocket(), fresh.herdrSocket())) { changed.add("herdrSocket"); } + if (!Objects.equals(old.memberHerdrSocket(), fresh.memberHerdrSocket())) { + changed.add("memberHerdrSocket"); + } if (!Objects.equals(old.broker(), fresh.broker())) { changed.add("broker"); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index a3183e9..96b3664 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -32,7 +32,8 @@ import java.util.Set; * silently dropping a whole block is indistinguishable from honouring it. * * @param bind REST/MCP listen host:port - * @param herdrSocket path to herdr's Unix socket ({@code null} → client default) + * @param herdrSocket path to the lead herdr Unix socket ({@code null} → client default) + * @param memberHerdrSocket optional member herdr Unix socket ({@code null}/blank → lead socket) * @param profiles named backend profiles, keyed by profile name (multi-backend fleet). A * profile answers which backend — model, CLI adapter, credentials, * cost. It says nothing about what the member spawned on it is for; that is @@ -80,6 +81,7 @@ import java.util.Set; public record FleetConfig( Bind bind, String herdrSocket, + String memberHerdrSocket, Map profiles, Guard guard, String worktreeRoot, @@ -105,7 +107,7 @@ public record FleetConfig( LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth, ConfigReload configReload, Integer quarantineCooldownSeconds, MemberCredentials memberCredentials) { - this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, + this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth, configReload, quarantineCooldownSeconds, memberCredentials, null); } @@ -116,9 +118,9 @@ public record FleetConfig( Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth, ConfigReload configReload, Integer quarantineCooldownSeconds) { - this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, + this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth, - configReload, quarantineCooldownSeconds, null); + configReload, quarantineCooldownSeconds, null, null); } /** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */ @@ -129,8 +131,8 @@ public record FleetConfig( String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs, Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, LeadHeartbeat leadHeartbeat, String placement, Auth auth) { - this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, - spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null); + this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, + spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null, null, null); } /** Back-compat form before the optional {@code health:} block was added. */ @@ -138,8 +140,8 @@ public record FleetConfig( String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs, Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) { - this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, - spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null); + this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, + spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null, null, null); } /** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */ @@ -148,9 +150,9 @@ public record FleetConfig( Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet, LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth, ConfigReload configReload) { - this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, + this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs, spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth, - configReload, null); + configReload, null, null, null); } /** @@ -1320,7 +1322,7 @@ public record FleetConfig( * {@code fleetd.yaml} itself is gitignored. */ static final Set KNOWN_TOP_LEVEL_KEYS = Set.of( - "bind", "herdrSocket", "profiles", "guard", "worktreeRoot", + "bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet", "leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds", "memberCredentials", "coordinator"); @@ -1941,7 +1943,7 @@ public record FleetConfig( : new MemberCredentials(null, List.of(), List.of()); // coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened, // and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup). - return new FleetConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs, + return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs, broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload, quarantineCooldown, mc, coordinator); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/herdr/HerdrRouter.java b/fleetd/src/main/java/dev/ltms/fleet/herdr/HerdrRouter.java new file mode 100644 index 0000000..29ffa56 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/herdr/HerdrRouter.java @@ -0,0 +1,40 @@ +package dev.ltms.fleet.herdr; + +import java.util.Objects; +import java.util.function.Predicate; + +/** Routes lead operations and member operations to their owning herdr daemon. */ +public final class HerdrRouter implements AutoCloseable { + private final HerdrClient lead; + private final HerdrClient member; + private final AgentControl leadAgents; + private final AgentControl memberAgents; + private final WorkspaceControl leadSpaces; + private final WorkspaceControl memberSpaces; + private final Predicate isLead; + + public HerdrRouter(HerdrClient lead, HerdrClient member, Predicate isLead) { + this.lead = Objects.requireNonNull(lead, "lead"); + this.member = member != null ? member : lead; + this.isLead = Objects.requireNonNull(isLead, "isLead"); + leadAgents = new AgentControl(this.lead); + memberAgents = this.member == this.lead ? leadAgents : new AgentControl(this.member); + leadSpaces = new WorkspaceControl(this.lead); + memberSpaces = this.member == this.lead ? leadSpaces : new WorkspaceControl(this.member); + } + + public AgentControl leadAgents() { return leadAgents; } + public WorkspaceControl leadSpaces() { return leadSpaces; } + public AgentControl memberAgents() { return memberAgents; } + public WorkspaceControl memberSpaces() { return memberSpaces; } + public AgentControl agentsFor(String targetId) { return isLead.test(targetId) ? leadAgents : memberAgents; } + + HerdrClient leadClient() { return lead; } + HerdrClient memberClient() { return member; } + + @Override + public void close() { + lead.close(); + if (member != lead) member.close(); + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java index 9bafae7..6dd05bd 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java @@ -2,6 +2,7 @@ package dev.ltms.fleet.inject; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentStatus; +import dev.ltms.fleet.herdr.HerdrRouter; import dev.ltms.fleet.msg.TurnToken; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -89,6 +90,7 @@ public final class Injector { public static final long POLL_INTERVAL_MILLIS = 250; private final AgentControl agents; + private final HerdrRouter router; private final TurnListener turnListener; private final Predicate ready; // CB-113: a target is deliverable only when available private final Consumer forget; // CB-114: clear a gone worker's readiness/presence @@ -124,11 +126,25 @@ public final class Injector { public Injector(AgentControl agents, TurnListener turnListener, Predicate ready, Consumer forget) { this.agents = agents; + this.router = null; this.turnListener = turnListener; this.ready = ready; this.forget = forget; } + public Injector(HerdrRouter router, TurnListener turnListener, Predicate ready, + Consumer forget) { + this.agents = null; + this.router = router; + this.turnListener = turnListener; + this.ready = ready; + this.forget = forget; + } + + private AgentControl agentsFor(String target) { + return router != null ? router.agentsFor(target) : agents; + } + /** A pending message and the future that completes when it has been delivered. */ private record Pending(String text, TurnToken token, CompletableFuture delivered) { } @@ -253,7 +269,7 @@ public final class Injector { if (p != null && ready.test(target)) { t.notReadySincePoll = 0; try { - agents.send(target, p.text()); + agentsFor(target).send(target, p.text()); t.queue.poll(); t.awaitingPickup = true; t.awaitingCompletion = true; @@ -313,7 +329,7 @@ public final class Injector { // thread while it holds the target lock. if (resubmit) { try { - agents.submit(target); // nudge a raced Enter so the pending paste submits + agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits } catch (RuntimeException e) { log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage()); } 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 706651f..87180ea 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusPoller.java @@ -3,6 +3,7 @@ package dev.ltms.fleet.inject; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentStatus; import dev.ltms.fleet.herdr.HerdrException; +import dev.ltms.fleet.herdr.HerdrRouter; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -22,6 +23,7 @@ public final class StatusPoller { private static final Logger log = LoggerFactory.getLogger(StatusPoller.class); private final AgentControl agents; + private final HerdrRouter router; private final Injector injector; private final StatusRefiner refiner; private final long intervalMillis; @@ -35,11 +37,20 @@ public final class StatusPoller { public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner, long intervalMillis) { this.agents = agents; + this.router = null; this.injector = injector; this.refiner = refiner; this.intervalMillis = intervalMillis; } + public StatusPoller(HerdrRouter router, Injector injector, long intervalMillis) { + this.agents = null; + this.router = router; + this.injector = injector; + this.refiner = null; + this.intervalMillis = intervalMillis; + } + /** Start the polling loop on a virtual thread. Idempotent. */ public synchronized void start() { if (running) return; @@ -56,7 +67,8 @@ public final class StatusPoller { try { // herdr's agent_status can misreport a settled worker as `unknown`; refine it // against the pane content before it drives delivery/completion (CB-115). - AgentStatus status = refiner.refine(target, agents.status(target)); + AgentControl control = router != null ? router.agentsFor(target) : agents; + AgentStatus status = router != null ? control.status(target) : 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. diff --git a/fleetd/src/test/java/dev/ltms/fleet/herdr/HerdrRouterTest.java b/fleetd/src/test/java/dev/ltms/fleet/herdr/HerdrRouterTest.java new file mode 100644 index 0000000..695bb2e --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/herdr/HerdrRouterTest.java @@ -0,0 +1,31 @@ +package dev.ltms.fleet.herdr; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertNotSame; + +class HerdrRouterTest { + @Test + void absentMemberClientSharesControlsForEveryTarget() { + FakeHerdr client = new FakeHerdr(); + HerdrRouter router = new HerdrRouter(client, null, id -> id.equals("lead")); + + assertSame(router.leadAgents(), router.memberAgents()); + assertSame(router.leadSpaces(), router.memberSpaces()); + assertSame(router.leadAgents(), router.agentsFor("lead")); + assertSame(router.leadAgents(), router.agentsFor("member")); + } + + @Test + void separateClientsRouteLeadAndMemberTargets() { + FakeHerdr lead = new FakeHerdr(); + FakeHerdr member = new FakeHerdr(); + HerdrRouter router = new HerdrRouter(lead, member, id -> id.equals("lead")); + + assertNotSame(router.leadAgents(), router.memberAgents()); + assertNotSame(router.leadSpaces(), router.memberSpaces()); + assertSame(router.leadAgents(), router.agentsFor("lead")); + assertSame(router.memberAgents(), router.agentsFor("member")); + } +}