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..c63aa07 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,12 @@ 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; + 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. @@ -169,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(agents, spaces, 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(agents, spaces, + adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(), opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(), () -> config.get().fleet(), @@ -271,6 +275,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)", @@ -278,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(agents, spaces, cfg).ensureLeads(); + int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads(); if (launched > 0) { log.info("lead auto-launch: {} lead(s) started", launched); } @@ -340,6 +346,7 @@ public final class Fleetd { + "BACKEND_EXHAUSTED): {}", credentialId, cfg.quarantineCooldownSeconds(), profile.profile(), reason); }); + 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 +388,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 @@ -420,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. @@ -430,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); @@ -439,7 +446,7 @@ public final class Fleetd { heartbeat = null; heartbeatScheduler.shutdownNow(); } - MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, + 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. @@ -491,8 +498,11 @@ 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. + // 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), new LsofPeerPidLookup(), new LsofProcessCwdLookup()); + 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. @@ -538,7 +548,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; @@ -589,10 +599,12 @@ 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(), + // 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).build(); app.start(cfg.bind().host(), cfg.bind().port()); log.info("fleetd listening on {}:{}, herdr socket {}", 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/herdr/PaneLocator.java b/fleetd/src/main/java/dev/ltms/fleet/herdr/PaneLocator.java index 9e2e355..f2a0af4 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/herdr/PaneLocator.java +++ b/fleetd/src/main/java/dev/ltms/fleet/herdr/PaneLocator.java @@ -2,6 +2,7 @@ package dev.ltms.fleet.herdr; import com.fasterxml.jackson.databind.JsonNode; +import java.util.List; import java.util.Map; /** @@ -13,33 +14,61 @@ import java.util.Map; *

herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid} * and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is * the obvious optimization once wired into {@code ClaudeCodeLauncher}. + * + *

CB-185 split the fleet across two herdr daemons — lead operations on one, members on the + * other ({@code memberHerdrSocket}). A caller's pane can live on either daemon (a lead's + * MCP connection resolves against the lead daemon; a member's against the member daemon), so this + * must be able to search more than one client. {@link #PaneLocator(HerdrClient, HerdrClient)} + * searches the lead client first, then the member client, and collapses to a single scan when the + * two are the same object (the historical single-daemon deployment). */ public final class PaneLocator { - private final HerdrClient herdr; + private final List herdrs; + /** Search only this client — the single-daemon deployment. */ public PaneLocator(HerdrClient herdr) { - this.herdr = herdr; + this.herdrs = List.of(herdr); + } + + /** + * Search {@code lead} first, then {@code member} — the two-daemon deployment (CB-185). When + * the caller passes the same client for both (no {@code memberHerdrSocket} configured), this + * collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour. + */ + public PaneLocator(HerdrClient lead, HerdrClient member) { + this.herdrs = lead == member ? List.of(lead) : List.of(lead, member); } /** * The {@code terminal_id} of the agent pane whose process tree contains {@code pid}, or - * {@code null} if no agent pane owns it (e.g. the caller is the primary, or off-host). + * {@code null} if no agent pane on any searched daemon owns it (e.g. the caller is the + * primary, or off-host). */ public String terminalForPid(long pid) { if (pid <= 0) { return null; } + for (HerdrClient herdr : herdrs) { + String terminal = terminalForPid(herdr, pid); + if (terminal != null) { + return terminal; + } + } + return null; + } + + private static String terminalForPid(HerdrClient herdr, long pid) { for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) { String paneId = pane.path("pane_id").asText(null); - if (paneId != null && paneOwnsPid(paneId, pid)) { + if (paneId != null && paneOwnsPid(herdr, paneId, pid)) { return pane.path("terminal_id").asText(null); } } return null; } - private boolean paneOwnsPid(String paneId, long pid) { + private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) { JsonNode info; try { info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info"); 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..a2e7966 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,24 @@ 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; + // CB-185: this refiner's own AgentControl (member) is only a default for the legacy 2-arg + // refine() overload — the loop below always calls the 3-arg refine(target, raw, control) + // with the per-target control from router.agentsFor(target), so a lead target is refined + // against the LEAD daemon even though this field points at the member one. + this.refiner = new StatusRefiner(router.memberAgents()); + this.intervalMillis = intervalMillis; + } + /** Start the polling loop on a virtual thread. Idempotent. */ public synchronized void start() { if (running) return; @@ -56,7 +71,11 @@ 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)); + // CB-185: refine THROUGH the same control the raw status came from — a router + // splits lead/member targets across two herdr daemons, and reading a lead's pane + // through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN. + AgentControl control = router != null ? router.agentsFor(target) : agents; + AgentStatus status = refiner.refine(target, control.status(target), control); injector.onStatus(target, status); } catch (HerdrException e) { // The worker's agent is gone — stop trying and unblock its waiters. diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusRefiner.java b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusRefiner.java index bc55d02..309b9b5 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/StatusRefiner.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/StatusRefiner.java @@ -41,15 +41,32 @@ public final class StatusRefiner { } /** - * Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned - * unchanged; an {@code UNKNOWN} triggers a pane read and content classification. A read failure - * leaves it {@code UNKNOWN} (the safe default: no delivery, and the stall path still applies). + * Return a trustworthy status for {@code target}, reading its pane through this refiner's own + * {@link AgentControl}. Equivalent to {@link #refine(String, AgentStatus, AgentControl)} with + * that control — kept for callers that only ever talk to one herdr daemon. */ public AgentStatus refine(String target, AgentStatus raw) { + return refine(target, raw, agents); + } + + /** + * Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned + * unchanged; an {@code UNKNOWN} triggers a pane read (through {@code control}) and content + * classification. A read failure leaves it {@code UNKNOWN} (the safe default: no delivery, and + * the stall path still applies). + * + *

CB-185: {@code control} must be the {@link AgentControl} for the same daemon the + * raw status was sampled from — a router splits lead and member targets across two herdr + * daemons, and reading a lead's pane through the member client (or vice versa) fails to find + * the pane and leaves the target wedged at {@code UNKNOWN} forever. Callers that route per + * target (e.g. {@code StatusPoller}) must pass that target's control explicitly rather than + * relying on the control fixed at construction. + */ + public AgentStatus refine(String target, AgentStatus raw, AgentControl control) { if (raw != AgentStatus.UNKNOWN) return raw; String pane; try { - pane = agents.read(target, PROBE_SOURCE); + pane = control.read(target, PROBE_SOURCE); } catch (RuntimeException e) { log.debug("status refine read for {} failed; leaving UNKNOWN: {}", target, e.getMessage()); return AgentStatus.UNKNOWN; diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 4cb6b19..6e7ff5e 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -2,6 +2,7 @@ package dev.ltms.fleet.msg; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentStatus; +import dev.ltms.fleet.herdr.HerdrRouter; import dev.ltms.fleet.inject.Injector; import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; @@ -189,6 +190,7 @@ public final class MessageService { } private final AgentControl agents; + private final HerdrRouter router; private final Injector injector; private final Rendezvous rendezvous; private final ReplyInbox inbox; @@ -257,6 +259,7 @@ public final class MessageService { MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) { this.agents = agents; + this.router = null; this.injector = injector; this.rendezvous = rendezvous; this.inbox = inbox; @@ -265,6 +268,20 @@ public final class MessageService { this.nowNanos = nowNanos; } + public MessageService(HerdrRouter router, Injector injector, Rendezvous rendezvous, ReplyInbox inbox, + ReplyPushLoop pushLoop, Metrics metrics) { + this.agents = null; + this.router = router; + this.injector = injector; + this.rendezvous = rendezvous; + this.inbox = inbox; + this.pushLoop = pushLoop; + this.metrics = metrics; + this.nowNanos = System::nanoTime; + } + + private AgentControl agentsFor(String target) { return router != null ? router.agentsFor(target) : agents; } + /** Create with an explicit {@link ReplyInbox} and no push loop. */ public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) { this(agents, injector, rendezvous, inbox, null); @@ -277,7 +294,7 @@ public final class MessageService { /** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */ public AgentStatus status(String target) { - return agents.status(target); + return agentsFor(target).status(target); } /** Read-only delegation fact for fleet views. */ @@ -788,7 +805,7 @@ public final class MessageService { /** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */ private String liveStatus(String target) { try { - return agents.status(target).name().toLowerCase(); + return agentsFor(target).status(target).name().toLowerCase(); } catch (RuntimeException e) { return "unknown"; } diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index e634390..16a0cf0 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -55,7 +55,8 @@ public final class FleetApp { /** Context attribute under which the resolved caller is stashed by the auth filter. */ private static final String CALLER = "fleetd.caller"; - private final HerdrClient herdr; + private final HerdrClient herdr; // lead daemon + private final HerdrClient memberHerdr; // CB-185: member daemon (same object when unconfigured) private final PeerLauncher workers; private final SessionManager sessions; // CB-301: authoritative session registry private final MessageService messages; @@ -95,7 +96,23 @@ public final class FleetApp { MessageService messages, MemberPresence presence, HttpServlet mcpServlet, CallerResolver auth, Metrics metrics, Predicate deliverable) { + this(herdr, herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable); + } + + /** + * @param herdr the lead daemon's client + * @param memberHerdr the member daemon's client (CB-185); pass the same instance as + * {@code herdr} for a single-daemon deployment — {@code healthz}/{@code + * sessions} then make exactly one herdr call each, unchanged from before + * the two-daemon router existed + * @param deliverable the injector's readiness gate, shared so status reports its real result + */ + public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions, + MessageService messages, MemberPresence presence, + HttpServlet mcpServlet, CallerResolver auth, Metrics metrics, + Predicate deliverable) { this.herdr = herdr; + this.memberHerdr = memberHerdr != null ? memberHerdr : herdr; this.workers = workers; this.sessions = sessions; this.messages = messages; @@ -190,30 +207,64 @@ public final class FleetApp { ctx.status(200).contentType("text/plain; version=0.0.4; charset=utf-8").result(metrics.render()); } - /** Liveness + herdr reachability. 200 when herdr answers ping, 503 otherwise. */ + /** + * Liveness + herdr reachability. 200 only when BOTH daemons answer ping — 503 otherwise + * (CB-185). With no {@code memberHerdrSocket} configured {@code memberHerdr == herdr}, so this + * makes exactly the one {@code ping} call it always did and reports the same body; with a + * second daemon configured, a member daemon that is down must not be masked by a healthy lead + * daemon — every spawn goes through the member daemon and would otherwise fail silently behind + * a green {@code /healthz}. + */ private void healthz(Context ctx) { + JsonNode pong; try { - JsonNode pong = herdr.call("ping"); - ctx.status(200).json(Map.of( - "status", "ok", - "herdr", Map.of( - "version", pong.path("version").asText(""), - "protocol", pong.path("protocol").asInt()))); + pong = herdr.call("ping"); } catch (HerdrException e) { ctx.status(503).json(Map.of( "status", "degraded", "herdr", "unreachable", "detail", e.getMessage())); + return; } + if (memberHerdr != herdr) { + try { + memberHerdr.call("ping"); + } catch (HerdrException e) { + ctx.status(503).json(Map.of( + "status", "degraded", + "herdr", "member unreachable", + "detail", e.getMessage())); + return; + } + } + ctx.status(200).json(Map.of( + "status", "ok", + "herdr", Map.of( + "version", pong.path("version").asText(""), + "protocol", pong.path("protocol").asInt()))); } - /** Sessions view derived from herdr {@code workspace.list} (one workspace → one row). */ + /** + * Sessions view derived from herdr {@code workspace.list} (one workspace → one row), merged + * across both daemons (CB-185). With no {@code memberHerdrSocket} configured {@code + * memberHerdr == herdr}, so this calls {@code workspace.list} exactly once, same as before the + * router existed; with a second daemon configured, calling it twice would silently drop every + * member workspace (they live on the member daemon only). + */ private void sessions(Context ctx) { if (!allow(ctx, Authz.Action.READ, null)) { return; } - JsonNode result = herdr.call("workspace.list"); List> out = new ArrayList<>(); + collectSessions(herdr, out); + if (memberHerdr != herdr) { + collectSessions(memberHerdr, out); + } + ctx.status(200).json(Map.of("sessions", out)); + } + + private static void collectSessions(HerdrClient client, List> out) { + JsonNode result = client.call("workspace.list"); for (JsonNode w : result.path("workspaces")) { out.add(Map.of( "id", w.path("workspace_id").asText(""), @@ -222,7 +273,6 @@ public final class FleetApp { "paneCount", w.path("pane_count").asInt(), "agentStatus", w.path("agent_status").asText("unknown"))); } - ctx.status(200).json(Map.of("sessions", out)); } /** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */ diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdConnectionIdentityConstructionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdConnectionIdentityConstructionTest.java new file mode 100644 index 0000000..7274264 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdConnectionIdentityConstructionTest.java @@ -0,0 +1,31 @@ +package dev.ltms.fleet; + +import java.nio.file.Files; +import java.nio.file.Path; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CB-185: {@code ConnectionIdentity} must resolve a caller's pane on EITHER herdr daemon (a + * lead's MCP connection resolves against the lead daemon; a member's against the member daemon). + * Pinning {@code PaneLocator} to {@code memberHerdr} alone — the bug this guards against — leaves + * every lead's own connection unresolvable ({@code callerTerminal == null}) the moment + * {@code memberHerdrSocket} names a second daemon, which breaks {@code fleet_reply}/{@code + * fleet_ask} and {@code fleet_whoami} for a lead. A unit test on {@link + * dev.ltms.fleet.herdr.PaneLocator} alone (see {@code PaneLocatorTest}) proves the class CAN + * search two clients, but not that {@code Fleetd.main} actually wires it that way — hence this + * source-level assertion, the same technique {@code FleetdHerdrControlConstructionTest} uses. + */ +class FleetdConnectionIdentityConstructionTest { + @Test + void connectionIdentitySearchesBothDaemonsNotJustTheMemberOne() throws Exception { + String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java")); + assertFalse(source.contains("new PaneLocator(memberHerdr)"), + "PaneLocator must not be pinned to the member daemon alone — a lead's own " + + "connection resolves against the LEAD daemon and would never be found"); + assertTrue(source.contains("new PaneLocator(herdr, memberHerdr)"), + "PaneLocator must search the lead daemon first, then the member daemon"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdFleetAppConstructionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdFleetAppConstructionTest.java new file mode 100644 index 0000000..a7e8482 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdFleetAppConstructionTest.java @@ -0,0 +1,29 @@ +package dev.ltms.fleet; + +import java.nio.file.Files; +import java.nio.file.Path; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CB-185: {@code FleetApp} must be constructed with BOTH herdr clients (the lead's and the + * member's), never the raw lead-only {@code herdr}. Passing only {@code herdr} — the bug this + * guards against — makes {@code GET /healthz} green while the member daemon is down (so every + * spawn fails invisibly) and silently drops every member workspace from {@code GET /sessions}. + * A behavioural test on {@code FleetApp} alone (see {@code FleetAppTwoDaemonTest}) proves the + * class merges/gates correctly when given two clients, but not that {@code Fleetd.main} actually + * passes it two — hence this source-level assertion, mirroring + * {@code FleetdHerdrControlConstructionTest}. + */ +class FleetdFleetAppConstructionTest { + @Test + void fleetAppIsConstructedWithBothHerdrDaemons() throws Exception { + String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java")); + assertFalse(source.contains("new FleetApp(herdr, workers,"), + "FleetApp must not be constructed with the lead-only herdr client"); + assertTrue(source.contains("new FleetApp(herdr, memberHerdr, workers,"), + "FleetApp must be constructed with both the lead and the member herdr client"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdHerdrControlConstructionTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdHerdrControlConstructionTest.java new file mode 100644 index 0000000..01dedf3 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdHerdrControlConstructionTest.java @@ -0,0 +1,17 @@ +package dev.ltms.fleet; + +import java.nio.file.Files; +import java.nio.file.Path; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; + +class FleetdHerdrControlConstructionTest { + @Test + void fleetdDelegatesStatefulControlsToTheRouter() throws Exception { + // AgentControl caches paneByTerminal, so the router must be its only production factory. + String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java")); + assertFalse(source.contains("new AgentControl(")); + assertFalse(source.contains("new WorkspaceControl(")); + } +} 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 cf86610..591b39c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/herdr/FakeHerdr.java +++ b/fleetd/src/test/java/dev/ltms/fleet/herdr/FakeHerdr.java @@ -40,6 +40,7 @@ public final class FakeHerdr implements HerdrClient { private int workerTabPaneCount = 1; private String paneCloseErrorCode = null; private String agentSendErrorCode = null; + private boolean noPanes = false; private volatile String agentStatus = "idle"; // steady-state agent.get status private volatile String readText = "worker transcript tail"; // canned agent.read output private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane @@ -75,6 +76,15 @@ public final class FakeHerdr implements HerdrClient { return this; } + /** + * Make {@code pane.list} report no panes at all — models a second herdr daemon (CB-185) that + * simply does not host the pane a {@link PaneLocator} is searching for. + */ + public FakeHerdr withNoPanes() { + this.noPanes = true; + return this; + } + /** Set the {@code agent_status} that {@code agent.get} reports (drives the injector). */ public FakeHerdr agentStatus(String status) { this.agentStatus = status; @@ -268,7 +278,9 @@ public final class FakeHerdr implements HerdrClient { case "pane.get" -> mapper.readTree(""" {"type":"pane_info","pane":{"pane_id":"w9:pW","workspace_id":"w9", "tab_id":"w9:t2","agent_status":"idle"}}"""); - case "pane.list" -> mapper.readTree(""" + case "pane.list" -> noPanes + ? mapper.readTree("{\"type\":\"pane_list\",\"panes\":[]}") + : mapper.readTree(""" {"type":"pane_list","panes":[ {"pane_id":"w2:p7","terminal_id":"term_a","workspace_id":"w2","tab_id":"w2:t7","agent":"claude"}, {"pane_id":"w2:p9","terminal_id":"term_shell","workspace_id":"w2","tab_id":"w2:t8"}]}"""); 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")); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/herdr/PaneLocatorTest.java b/fleetd/src/test/java/dev/ltms/fleet/herdr/PaneLocatorTest.java index 7b01877..929ddc1 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/herdr/PaneLocatorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/herdr/PaneLocatorTest.java @@ -24,4 +24,43 @@ class PaneLocatorTest { assertNull(loc.terminalForPid(0)); assertNull(loc.terminalForPid(-1)); } + + // --- two-daemon fallback (CB-185) ----------------------------------------- + + @Test + void fallsBackToTheMemberClientWhenTheLeadHasNoMatch() { + // The caller's pane lives on the member daemon only (e.g. the caller is a spawned + // member) — the lead client reports no panes at all, so the locator must fall back. + HerdrClient lead = new FakeHerdr().withNoPanes(); + HerdrClient member = new FakeHerdr(); + PaneLocator two = new PaneLocator(lead, member); + assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID)); + } + + @Test + void searchesTheLeadClientBeforeTheMemberClient() { + // The caller's pane lives on the LEAD daemon (e.g. the caller is a peer lead) — with two + // daemons, resolving it must not depend on the member client having a matching pane. + HerdrClient lead = new FakeHerdr(); + HerdrClient member = new FakeHerdr().withNoPanes(); + PaneLocator two = new PaneLocator(lead, member); + assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID)); + } + + @Test + void nullWhenNeitherClientHasTheMatch() { + PaneLocator two = new PaneLocator(new FakeHerdr().withNoPanes(), new FakeHerdr().withNoPanes()); + assertNull(two.terminalForPid(FakeHerdr.WORKER_PID)); + } + + @Test + void collapsesToOneScanWhenLeadAndMemberAreTheSameClient() { + // The single-daemon deployment (no memberHerdrSocket configured): the two-arg constructor + // must behave exactly like the one-arg constructor, including making only one herdr call. + FakeHerdr shared = new FakeHerdr(); + PaneLocator two = new PaneLocator(shared, shared); + assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID)); + long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count(); + assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice"); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java new file mode 100644 index 0000000..d748534 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusPollerRoutingTest.java @@ -0,0 +1,75 @@ +package dev.ltms.fleet.inject; + +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrRouter; +import dev.ltms.fleet.msg.TestTurnTokens; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CB-185: with a router split across two herdr daemons, {@link StatusPoller} must refine a raw + * {@code UNKNOWN} status by reading the pane content from the SAME daemon the status was sampled + * from — the lead daemon for a lead target, the member daemon for a member target. Reading the + * wrong daemon never finds the pane, classification stays {@code UNKNOWN} forever, and the + * status-gated {@link Injector} wedges: a queued message is never delivered. + * + *

This exercises the real production classes ({@code StatusPoller(HerdrRouter, ...)}, + * {@code Injector(HerdrRouter, ...)}) wired together, not a hand-built object graph — the earlier + * three CB-185 bugs all passed exactly that kind of test while the real wiring stayed broken. + */ +class StatusPollerRoutingTest { + + private static final String LEAD_TARGET = "term_a"; + + @Test + void refinesALeadTargetFromTheLeadDaemonAndDelivers() throws Exception { + // The lead daemon's pane is at a settled idle prompt; the member daemon's pane content is + // unclassifiable garbage. A correct refiner reads the LEAD daemon and delivers. + FakeHerdr leadHerdr = new FakeHerdr().agentStatus("unknown").readText("⏺ answer\n❯ "); + FakeHerdr memberHerdr = new FakeHerdr().agentStatus("unknown") + .readText("garbled ansi noise with no prompt"); + HerdrRouter router = new HerdrRouter(leadHerdr, memberHerdr, LEAD_TARGET::equals); + Injector injector = new Injector(router, TurnListener.NOOP, _ -> true, _ -> { + }); + StatusPoller poller = new StatusPoller(router, injector, 10); + poller.start(); + try { + CompletableFuture delivered = + injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)); + // Must resolve quickly: refining against the WRONG daemon (member) never classifies + // out of UNKNOWN, so this would time out under the bug. + delivered.get(2, TimeUnit.SECONDS); + } finally { + poller.stop(); + } + assertTrue(leadHerdr.called("agent.read"), "refine must probe the LEAD daemon's pane content"); + } + + @Test + void aLeadTargetNeverDeliversWhenOnlyTheMemberDaemonIsClassifiable() throws Exception { + // Inverted control: the member daemon's content WOULD classify to idle, but this is a lead + // target — a correct implementation must not use it, so delivery must NOT happen. + FakeHerdr leadHerdr = new FakeHerdr().agentStatus("unknown") + .readText("garbled ansi noise with no prompt"); + FakeHerdr memberHerdr = new FakeHerdr().agentStatus("unknown").readText("⏺ answer\n❯ "); + HerdrRouter router = new HerdrRouter(leadHerdr, memberHerdr, LEAD_TARGET::equals); + Injector injector = new Injector(router, TurnListener.NOOP, _ -> true, _ -> { + }); + StatusPoller poller = new StatusPoller(router, injector, 10); + poller.start(); + try { + CompletableFuture delivered = + injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)); + assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS), + "a lead target must never be refined from the member daemon's pane content"); + } finally { + poller.stop(); + } + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusRefinerTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusRefinerTest.java index 7034bea..642eadd 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/StatusRefinerTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/StatusRefinerTest.java @@ -85,4 +85,34 @@ class StatusRefinerTest { assertEquals(AgentStatus.UNKNOWN, refiner.refine("term_a", AgentStatus.UNKNOWN)); } + + // --- refine(target, raw, control) — CB-185 per-call routing --------------- + + @Test + void threeArgRefineReadsThroughTheGivenControlNotTheConstructedOne() { + // The refiner is CONSTRUCTED with one control (standing in for "the member daemon"), but + // a call names a DIFFERENT control (standing in for "the lead daemon") — the read must go + // to the one passed to the call, since that is the daemon the raw status came from. + FakeHerdr constructedWith = new FakeHerdr().readText("nothing recognizable here"); + FakeHerdr passedToCall = new FakeHerdr().readText("⏺ answer\n❯ "); + StatusRefiner refiner = new StatusRefiner(new AgentControl(constructedWith)); + + AgentStatus result = refiner.refine("term_a", AgentStatus.UNKNOWN, new AgentControl(passedToCall)); + + assertEquals(AgentStatus.IDLE, result, "must classify from the PASSED control's pane content"); + assertTrue(passedToCall.called("agent.read")); + assertFalse(constructedWith.called("agent.read"), + "the control fixed at construction must not be read when a call-site control is given"); + } + + @Test + void twoArgRefineStillReadsTheConstructedControl() { + // The legacy 2-arg overload (single-daemon callers) must keep using the constructed + // control — this is refine(target, raw, control) called with the field as `control`. + FakeHerdr herdr = new FakeHerdr().readText("⏺ answer\n❯ "); + StatusRefiner refiner = new StatusRefiner(new AgentControl(herdr)); + + assertEquals(AgentStatus.IDLE, refiner.refine("term_a", AgentStatus.UNKNOWN)); + assertTrue(herdr.called("agent.read")); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTwoDaemonTest.java b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTwoDaemonTest.java new file mode 100644 index 0000000..efbba3d --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/rest/FleetAppTwoDaemonTest.java @@ -0,0 +1,99 @@ +package dev.ltms.fleet.rest; + +import dev.ltms.fleet.herdr.FakeHerdr; +import dev.ltms.fleet.herdr.HerdrClient; +import io.javalin.Javalin; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-185: with a router split across two herdr daemons (lead + {@code memberHerdrSocket}), + * {@link FleetApp#healthz} must require BOTH daemons to answer and {@link FleetApp#sessions} + * (which the {@code GET /sessions} route calls) must merge workspaces from both — the bug this + * guards against had {@code FleetApp} constructed with the raw lead-only client, so a down member + * daemon was invisible behind a green {@code /healthz} (every spawn then fails) and every member + * workspace was silently dropped from {@code GET /sessions}. + * + *

Builds the real {@link FleetApp} directly (not a hand-rolled stand-in) against only the two + * herdr clients — the other collaborators are unused by the two routes under test here. + */ +class FleetAppTwoDaemonTest { + + private final HttpClient http = HttpClient.newHttpClient(); + private Javalin app; + + @AfterEach + void stop() { + if (app != null) app.stop(); + } + + private int start(HerdrClient lead, HerdrClient member) { + app = new FleetApp(lead, member, null, null, null, null, null, null, null, ignored -> false) + .build().start("127.0.0.1", 0); + return app.port(); + } + + private HttpResponse get(int port, String path) throws Exception { + HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path)).GET().build(); + return http.send(req, HttpResponse.BodyHandlers.ofString()); + } + + @Test + void healthzIsGreenWhenBothDaemonsAnswer() throws Exception { + int port = start(new FakeHerdr(), new FakeHerdr()); + assertEquals(200, get(port, "/healthz").statusCode()); + } + + @Test + void healthzIsDegradedWhenOnlyTheMemberDaemonIsDown() throws Exception { + int port = start(new FakeHerdr(), new FakeHerdr().healthy(false)); + HttpResponse res = get(port, "/healthz"); + assertEquals(503, res.statusCode(), + "a down MEMBER daemon must not be masked by a healthy lead — every spawn goes " + + "through the member daemon"); + } + + @Test + void healthzIsDegradedWhenOnlyTheLeadDaemonIsDown() throws Exception { + int port = start(new FakeHerdr().healthy(false), new FakeHerdr()); + assertEquals(503, get(port, "/healthz").statusCode()); + } + + @Test + void healthzMakesExactlyOneCallWhenLeadAndMemberAreTheSameClient() throws Exception { + // Single-daemon deployment (no memberHerdrSocket) — must be byte-for-byte the old + // behaviour: one ping call, 200 on success. + FakeHerdr shared = new FakeHerdr(); + int port = start(shared, shared); + assertEquals(200, get(port, "/healthz").statusCode()); + long pings = shared.calls.stream().filter(c -> c.method().equals("ping")).count(); + assertEquals(1, pings, "single-daemon deployment must make exactly one ping call"); + } + + @Test + void sessionsMergesWorkspacesFromBothDaemons() throws Exception { + FakeHerdr lead = new FakeHerdr(); + FakeHerdr member = new FakeHerdr().withWorkspace("w9", "member-only-workspace"); + int port = start(lead, member); + HttpResponse res = get(port, "/sessions"); + assertEquals(200, res.statusCode(), res.body()); + assertTrue(res.body().contains("member-only-workspace"), + "GET /sessions must not silently drop the member daemon's workspaces"); + } + + @Test + void sessionsMakesExactlyOneWorkspaceListCallWhenLeadAndMemberAreTheSameClient() throws Exception { + FakeHerdr shared = new FakeHerdr(); + int port = start(shared, shared); + assertEquals(200, get(port, "/sessions").statusCode()); + long calls = shared.calls.stream().filter(c -> c.method().equals("workspace.list")).count(); + assertEquals(1, calls, "single-daemon deployment must call workspace.list exactly once"); + } +}