CB-185: fix two blockers to switching on memberHerdrSocket #196
@@ -3,6 +3,7 @@ package dev.ltms.fleet.member;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
@@ -46,9 +47,13 @@ import java.util.stream.Collectors;
|
||||
* the single adapter that declares it. Profiles partition cleanly across adapters: the
|
||||
* constructor rejects a name claimed by two.</li>
|
||||
* <li><strong>By pane id</strong> — {@link #stop} routes to the adapter that spawned that pane
|
||||
* (recorded at spawn time). A pane the composite never spawned can use the fallback route
|
||||
* in a one-daemon fleet. With more than one herdr daemon, its owner is unknown, so stop refuses
|
||||
* the ambiguous id rather than closing a pane on an arbitrary herdr daemon.</li>
|
||||
* (recorded at spawn time). A pane the composite never spawned, or one whose record was lost
|
||||
* to a daemon restart (CB-185 blocker 1 — {@link #spawnedBy} is in-memory only), can use the
|
||||
* fallback route in a one-daemon fleet. With more than one herdr daemon, {@link #probeOwner}
|
||||
* asks each configured daemon which one actually knows the pane: exactly one match routes
|
||||
* (and caches); no match is treated as already-gone; more than one match is a genuine
|
||||
* ambiguity (pane ids are per-daemon counters, so two daemons really can both hold, say,
|
||||
* {@code w1:p1}) and stop refuses rather than closing a pane on an arbitrary herdr daemon.</li>
|
||||
* <li><strong>Fleet-wide</strong> — {@link #reapOrphanWorkers} and {@link #capabilities} fan out
|
||||
* and combine. {@link #list} is deduplicated by (owning daemon, pane id): delegates that share
|
||||
* one herdr connection report the same global agent set, but two daemons can each hold a pane
|
||||
@@ -440,12 +445,23 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
public void stop(String id) {
|
||||
HerdrPeerLauncher d = spawnedBy.get(id);
|
||||
if (d == null) {
|
||||
if (herdrDaemonCount() != 1) {
|
||||
throw new IllegalArgumentException("ambiguous paneId '" + id
|
||||
+ "': no owning herdr daemon was recorded");
|
||||
if (herdrDaemonCount() == 1) {
|
||||
log.debug("stop({}) — no recorded owner in a single-daemon fleet", id);
|
||||
d = delegates.getFirst();
|
||||
} else {
|
||||
d = probeOwner(id);
|
||||
if (d == null) {
|
||||
// No configured herdr daemon has ever heard of this pane. CB-185 blocker 1: this
|
||||
// is the normal case right after a daemon restart empties spawnedBy for a member
|
||||
// that has ALREADY been torn down since — the caller retried a stop that already
|
||||
// succeeded. Nothing to close and no owner to cache; matching the tolerance
|
||||
// HerdrPeerLauncher#stop already gives an already-gone pane (agent.close swallows
|
||||
// that as success), stop() here is a no-op rather than a refusal.
|
||||
log.debug("stop({}) — no configured herdr daemon knows this pane; "
|
||||
+ "treating as already stopped", id);
|
||||
return;
|
||||
}
|
||||
}
|
||||
log.debug("stop({}) — no recorded owner in a single-daemon fleet", id);
|
||||
d = delegates.getFirst();
|
||||
}
|
||||
// Drop the owner record only after the delegate accepted the stop. Removing it first meant a
|
||||
// delegate that threw left the pane alive with its owner forgotten, so the retry fell into
|
||||
@@ -454,6 +470,61 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
spawnedBy.remove(id);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-185 blocker 1: recover a spawnedBy cache miss by asking every distinct herdr daemon which
|
||||
* one actually knows {@code id} — the fix for "after a restart, every surviving member becomes
|
||||
* un-stoppable" (spawnedBy is in-memory only, so a restart empties it, and members intentionally
|
||||
* outlive the daemon).
|
||||
*
|
||||
* <p>Grouped by daemon identity, not by delegate, for the same reason {@link #list()} groups
|
||||
* that way: two adapters (claude-code, opencode) sharing one herdr connection would otherwise be
|
||||
* probed twice, and a pane on their shared daemon would look owned by two adapters instead of
|
||||
* one daemon.
|
||||
*
|
||||
* <p>A daemon that fails to answer {@code list()} (e.g. it is down) is treated as "does not know
|
||||
* this pane" rather than aborting the whole probe — one unreachable daemon must never make a
|
||||
* pane that a <em>different</em>, healthy daemon actually owns un-stoppable too, which would
|
||||
* resurrect the exact bug this method exists to fix.
|
||||
*
|
||||
* @return the owning delegate — cached into {@link #spawnedBy} so the next call is free — or
|
||||
* {@code null} when no daemon knows the pane
|
||||
* @throws IllegalArgumentException when more than one daemon claims the pane: pane ids are
|
||||
* per-daemon counters, so two daemons really can both hold, say, {@code w1:p1}, and there
|
||||
* is no way to tell which one the caller means
|
||||
*/
|
||||
private HerdrPeerLauncher probeOwner(String id) {
|
||||
Map<HerdrClient, HerdrPeerLauncher> byDaemon = new IdentityHashMap<>();
|
||||
for (HerdrPeerLauncher delegate : delegates) {
|
||||
byDaemon.putIfAbsent(delegate.herdr(), delegate);
|
||||
}
|
||||
List<HerdrPeerLauncher> owners = new ArrayList<>();
|
||||
for (HerdrPeerLauncher representative : byDaemon.values()) {
|
||||
List<Agent> agents;
|
||||
try {
|
||||
agents = representative.list();
|
||||
} catch (HerdrException e) {
|
||||
log.warn("stop({}) probe: a configured herdr daemon was unreachable ({}); "
|
||||
+ "treating it as not knowing this pane", id, e.getClass().getSimpleName());
|
||||
continue;
|
||||
}
|
||||
boolean knows = agents.stream().anyMatch(a -> id.equals(a.paneId()));
|
||||
if (knows) {
|
||||
owners.add(representative);
|
||||
}
|
||||
}
|
||||
if (owners.size() > 1) {
|
||||
throw new IllegalArgumentException("ambiguous paneId '" + id + "': "
|
||||
+ owners.size() + " configured herdr daemons report this pane — "
|
||||
+ "no way to tell which one the caller means");
|
||||
}
|
||||
if (owners.isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
HerdrPeerLauncher owner = owners.get(0);
|
||||
spawnedBy.put(id, owner);
|
||||
return owner;
|
||||
}
|
||||
|
||||
/**
|
||||
* Count actual herdr daemons, not peer adapter kinds. Identity is intentional: separate client
|
||||
* objects may represent different daemons even if a client later implements value equality.
|
||||
|
||||
@@ -214,6 +214,17 @@ public final class FleetApp {
|
||||
* 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}.
|
||||
*
|
||||
* <p>CB-185 blocker 2: the {@code herdr} key always carries the <em>lead</em> daemon's
|
||||
* version/protocol, unchanged, because two consumers — {@code scripts/redeploy-fleetd.sh} and
|
||||
* {@code scripts/rename-checkout.sh} — read this endpoint already (both only check the HTTP
|
||||
* status code and print the body verbatim; neither parses a specific field, so adding a key
|
||||
* alongside {@code herdr} is safe). But it is the <em>member</em> daemon's protocol that decides
|
||||
* whether a spawn works, so when a second daemon is configured its version/protocol is reported
|
||||
* too, under a separate {@code member} key — never folded into {@code herdr}, which would make a
|
||||
* mismatch invisible to whichever consumer only reads that key. If the two protocol numbers
|
||||
* differ, {@code protocolMismatch: true} calls it out explicitly rather than leaving it to be
|
||||
* spotted by comparing two numbers by eye.
|
||||
*/
|
||||
private void healthz(Context ctx) {
|
||||
JsonNode pong;
|
||||
@@ -226,9 +237,15 @@ public final class FleetApp {
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("status", "ok");
|
||||
body.put("herdr", Map.of(
|
||||
"version", pong.path("version").asText(""),
|
||||
"protocol", pong.path("protocol").asInt()));
|
||||
if (memberHerdr != herdr) {
|
||||
JsonNode memberPong;
|
||||
try {
|
||||
memberHerdr.call("ping");
|
||||
memberPong = memberHerdr.call("ping");
|
||||
} catch (HerdrException e) {
|
||||
ctx.status(503).json(Map.of(
|
||||
"status", "degraded",
|
||||
@@ -236,12 +253,16 @@ public final class FleetApp {
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
int leadProtocol = pong.path("protocol").asInt();
|
||||
int memberProtocol = memberPong.path("protocol").asInt();
|
||||
body.put("member", Map.of(
|
||||
"version", memberPong.path("version").asText(""),
|
||||
"protocol", memberProtocol));
|
||||
if (leadProtocol != memberProtocol) {
|
||||
body.put("protocolMismatch", true);
|
||||
}
|
||||
}
|
||||
ctx.status(200).json(Map.of(
|
||||
"status", "ok",
|
||||
"herdr", Map.of(
|
||||
"version", pong.path("version").asText(""),
|
||||
"protocol", pong.path("protocol").asInt())));
|
||||
ctx.status(200).json(body);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -31,6 +31,8 @@ public final class FakeHerdr implements HerdrClient {
|
||||
*/
|
||||
public final List<Call> calls = new CopyOnWriteArrayList<>();
|
||||
private boolean healthy = true;
|
||||
private String pingVersion = "0.8.0";
|
||||
private int pingProtocol = 19;
|
||||
private final List<String> extraWorkspaces = new ArrayList<>();
|
||||
private final List<String> extraAgents = new ArrayList<>();
|
||||
/** workspaceId → extra tabs that {@code tab.list} reports for it (CB-558 lead scans). */
|
||||
@@ -52,6 +54,16 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code ping} report this version/protocol instead of the default 0.8.0/19 — CB-185
|
||||
* blocker 2's fixture for a lead and a member daemon running mismatched herdr versions.
|
||||
*/
|
||||
public FakeHerdr pingReports(String version, int protocol) {
|
||||
this.pingVersion = version;
|
||||
this.pingProtocol = protocol;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Reject the first {@code n} {@code agent.start} calls with {@code agent_name_taken}. */
|
||||
public FakeHerdr agentNameTakenTimes(int n) {
|
||||
this.agentNameTakenFor = n;
|
||||
@@ -166,7 +178,8 @@ public final class FakeHerdr implements HerdrClient {
|
||||
try {
|
||||
return switch (method) {
|
||||
case "ping" -> mapper.readTree(
|
||||
"{\"type\":\"pong\",\"version\":\"0.8.0\",\"protocol\":19}");
|
||||
("{\"type\":\"pong\",\"version\":\"%s\",\"protocol\":%d}")
|
||||
.formatted(pingVersion, pingProtocol));
|
||||
case "workspace.list" -> mapper.readTree(("""
|
||||
{"type":"workspace_list","workspaces":[
|
||||
{"workspace_id":"w1","label":"dev-mgnl","focused":true,"pane_count":7,"agent_status":"unknown"},
|
||||
|
||||
@@ -346,13 +346,94 @@ class CompositePeerLauncherTest {
|
||||
|
||||
@Test
|
||||
void stopRejectsAnUnownedPaneIdWhenMultipleDaemonsCouldOwnIt() {
|
||||
// CB-185 blocker 1: genuine ambiguity — pane ids are per-daemon counters, so two daemons
|
||||
// can each really hold an agent at "w1:p1". Neither claims ownership through spawnedBy
|
||||
// (empty, as after a restart), so the probe must find BOTH and refuse rather than guess.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
FakeHerdr second = new FakeHerdr().withAgent("y", "term_y", "w1:p1", "w1:t1");
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(new FakeHerdr()), opencodeAdapter(new FakeHerdr())), "claude");
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
IllegalArgumentException error = assertThrows(IllegalArgumentException.class,
|
||||
() -> composite.stop("w1:p1"));
|
||||
|
||||
assertEquals("ambiguous paneId 'w1:p1': no owning herdr daemon was recorded", error.getMessage());
|
||||
assertEquals("ambiguous paneId 'w1:p1': 2 configured herdr daemons report this pane — "
|
||||
+ "no way to tell which one the caller means", error.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopOnAPaneNoConfiguredDaemonKnowsIsTreatedAsAlreadyStopped() {
|
||||
// CB-185 blocker 1, the zero-owner branch: spawnedBy is empty (as after a restart) and
|
||||
// neither daemon's agent.list mentions this pane at all — it is already gone. A retried
|
||||
// stop() on an already-gone pane must succeed quietly, not refuse forever.
|
||||
FakeHerdr first = new FakeHerdr();
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
assertDoesNotThrow(() -> composite.stop("w1:p1"));
|
||||
|
||||
assertFalse(first.called("pane.close"), "no owner was found, so no delegate is told to close anything");
|
||||
assertFalse(second.called("pane.close"), "no owner was found, so no delegate is told to close anything");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopWithEmptySpawnedByResolvesTheOwnerThroughAProbeAndSkipsTheOtherDaemon() {
|
||||
// CB-185 blocker 1, the main fix: after a restart spawnedBy is empty for every surviving
|
||||
// member. stop() must still find the one daemon that actually knows the pane and route
|
||||
// only to it — never touching the daemon that never held it.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
composite.stop("w1:p1");
|
||||
|
||||
assertTrue(first.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w1:p1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the daemon that actually knows the pane closes it");
|
||||
assertFalse(second.called("pane.close"), "the daemon that never held the pane is never touched");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aProbeSurvivesOneUnreachableDaemonAndStillFindsTheOwnerOnTheOtherOne() {
|
||||
// CB-185 blocker 1 (lead review): a daemon that is DOWN while we probe must not abort the
|
||||
// whole probe — the pane the OPERATOR actually wants stopped can live on a different,
|
||||
// healthy daemon, and that pane must not become un-stoppable because a third one is down.
|
||||
FakeHerdr down = new FakeHerdr().healthy(false);
|
||||
FakeHerdr owner = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(down), opencodeAdapter(owner)), "claude");
|
||||
|
||||
assertDoesNotThrow(() -> composite.stop("w1:p1"),
|
||||
"the unreachable daemon must be skipped, not fail the whole stop");
|
||||
|
||||
assertTrue(owner.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w1:p1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the healthy daemon that actually owns the pane still closes it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aProbedOwnerIsCachedSoARetryAfterAFailedStopNeedsNoSecondProbe() {
|
||||
// CB-185 blocker 1: the probe's whole point is to be cheap on repeat — a failed stop (e.g.
|
||||
// "pane_busy") must not force another agent.list() round trip on every retry.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1")
|
||||
.paneCloseFailsWith("pane_busy");
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
assertThrows(HerdrException.class, () -> composite.stop("w1:p1"));
|
||||
long listCallsAfterFirst = first.calls.stream().filter(c -> c.method().equals("agent.list")).count()
|
||||
+ second.calls.stream().filter(c -> c.method().equals("agent.list")).count();
|
||||
assertTrue(listCallsAfterFirst > 0, "the first stop needed a probe");
|
||||
|
||||
assertThrows(HerdrException.class, () -> composite.stop("w1:p1"),
|
||||
"still failing on the retry, but through the cached owner");
|
||||
long listCallsAfterSecond = first.calls.stream().filter(c -> c.method().equals("agent.list")).count()
|
||||
+ second.calls.stream().filter(c -> c.method().equals("agent.list")).count();
|
||||
assertEquals(listCallsAfterFirst, listCallsAfterSecond,
|
||||
"the retry is served from the cache — no additional agent.list probe");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -96,4 +96,55 @@ class FleetAppTwoDaemonTest {
|
||||
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");
|
||||
}
|
||||
|
||||
// ── CB-185 blocker 2: /healthz must report the MEMBER daemon's protocol too ────────────────
|
||||
|
||||
@Test
|
||||
void healthzReportsBothDaemonsWhenTheirProtocolsDiffer() throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr().pingReports("0.8.0", 19);
|
||||
FakeHerdr member = new FakeHerdr().pingReports("0.7.0", 18);
|
||||
int port = start(lead, member);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
assertTrue(res.body().contains("\"protocol\":19"),
|
||||
"the herdr key keeps reporting the LEAD's protocol, unchanged: " + res.body());
|
||||
assertTrue(res.body().contains("\"member\""), "a separate member key is present: " + res.body());
|
||||
assertTrue(res.body().contains("\"protocol\":18"),
|
||||
"the member key reports the member daemon's own protocol: " + res.body());
|
||||
assertTrue(res.body().contains("\"protocolMismatch\":true"),
|
||||
"a differing protocol is called out explicitly, not left to be spotted by eye: " + res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzReportsBothDaemonsWithNoMismatchWhenProtocolsMatch() throws Exception {
|
||||
int port = start(new FakeHerdr(), new FakeHerdr());
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
assertTrue(res.body().contains("\"member\""), "the member key is present whenever a second daemon "
|
||||
+ "is configured, even when the protocols happen to agree: " + res.body());
|
||||
assertFalse(res.body().contains("protocolMismatch"),
|
||||
"matching protocols must not raise a mismatch flag: " + res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzWithOneDaemonCarriesNoMemberOrMismatchKey() throws Exception {
|
||||
// The single-daemon deployment (no memberHerdrSocket) must see no change at all beyond the
|
||||
// historical body: no "member" key, no "protocolMismatch" key. (Map.of()'s own key order is
|
||||
// JVM-salted regardless of this fix, so this checks content, not exact key order.)
|
||||
FakeHerdr shared = new FakeHerdr();
|
||||
int port = start(shared, shared);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
|
||||
assertEquals(200, res.statusCode());
|
||||
assertTrue(res.body().contains("\"status\":\"ok\""), res.body());
|
||||
assertTrue(res.body().contains("\"protocol\":19"), res.body());
|
||||
assertTrue(res.body().contains("\"version\":\"0.8.0\""), res.body());
|
||||
assertFalse(res.body().contains("\"member\""), "no second daemon configured, so no member key: " + res.body());
|
||||
assertFalse(res.body().contains("protocolMismatch"), res.body());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user