CB-185: route members to separate herdr #186

Merged
ltms merged 5 commits from worker/cb185-router-d6436d-3 into main 2026-08-29 01:24:32 +02:00
20 changed files with 625 additions and 54 deletions
+3
View File
@@ -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
+27 -15
View File
@@ -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<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.
@@ -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<String> 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 {}",
@@ -71,7 +71,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
/** Keys that cannot change under a running daemon — see the class doc. */
private static final Set<String> COLD_KEYS =
Set.of("bind", "herdrSocket", "broker", "auth");
Set.of("bind", "herdrSocket", "memberHerdrSocket", "broker", "auth");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -190,6 +190,9 @@ public final class ConfigRef implements Supplier<FleetConfig> {
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");
}
@@ -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 <em>which backend</em> — 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<String, Profile> 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<String> 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);
}
@@ -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<String> isLead;
public HerdrRouter(HerdrClient lead, HerdrClient member, Predicate<String> 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();
}
}
@@ -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;
* <p>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}.
*
* <p>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 <em>either</em> 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<HerdrClient> 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");
@@ -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<String> ready; // CB-113: a target is deliverable only when available
private final Consumer<String> 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<String> ready,
Consumer<String> forget) {
this.agents = agents;
this.router = null;
this.turnListener = turnListener;
this.ready = ready;
this.forget = forget;
}
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
Consumer<String> 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<Void> 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());
}
@@ -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.
@@ -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).
*
* <p>CB-185: {@code control} must be the {@link AgentControl} for the <em>same</em> 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;
@@ -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";
}
@@ -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<String> 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<String> 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<Map<String, Object>> 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<Map<String, Object>> 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. */
@@ -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");
}
}
@@ -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");
}
}
@@ -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("));
}
}
@@ -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"}]}""");
@@ -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"));
}
}
@@ -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");
}
}
@@ -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.
*
* <p>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<Void> 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<Void> 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();
}
}
}
@@ -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"));
}
}
@@ -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}.
*
* <p>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<String> 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<String> 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<String> 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");
}
}