CB-573: add fleet capacity view
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Successful in 1m25s

This commit is contained in:
Dai Ha
2026-08-15 05:43:46 +02:00
parent bf0ff2adbf
commit 01fab15713
4 changed files with 94 additions and 9 deletions
@@ -383,7 +383,11 @@ public final class Bridged {
}
BridgeMcp mcp = new BridgeMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics);
primaryRegistry, callers, metrics, profile -> liveCountRef.get().apply(profile),
profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.maxLoad();
}, System::nanoTime);
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
@@ -35,6 +35,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.stream.Collectors;
/**
@@ -74,6 +75,9 @@ public final class BridgeMcp {
private final McpSyncServer server;
private final CallerResolver authz; // CB-501: null → authorization not enforced (legacy)
private final Metrics metrics; // CB-502: null → auth failures not counted
private final Function<String, Integer> liveCount;
private final Function<String, Integer> maxLoad;
private final LongSupplier clock;
/**
* Legacy constructor — no authorization. Retained so existing tests exercise tool behaviour
@@ -82,7 +86,7 @@ public final class BridgeMcp {
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, MemberPresence presence,
PrimaryRegistry primaryRegistry) {
this(messages, workers, sessions, identity, presence, primaryRegistry, null, null);
this(messages, workers, sessions, identity, presence, primaryRegistry, null, null, _ -> 0, _ -> null, System::nanoTime);
}
/**
@@ -94,7 +98,18 @@ public final class BridgeMcp {
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, MemberPresence presence,
PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics) {
PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics,
_ -> 0, _ -> null, System::nanoTime);
}
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, Function<String, Integer> liveCount,
Function<String, Integer> maxLoad, LongSupplier clock) {
this.liveCount = liveCount;
this.maxLoad = maxLoad;
this.clock = clock;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
.jsonMapper(json)
@@ -208,7 +223,7 @@ public final class BridgeMcp {
.toolCall(listTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return listFleet(workers, sessions,
return listFleet(workers, sessions, messages, liveCount, maxLoad, clock,
callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange));
})
@@ -719,6 +734,12 @@ public final class BridgeMcp {
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, null, _ -> 0, _ -> null, System::nanoTime, leads, selfTerm);
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
LongSupplier clock, Map<String, String> leads, String selfTerm) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -728,15 +749,52 @@ public final class BridgeMcp {
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
.toList();
List<Map<String, Object>> out = sessions.roster().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
List<MemberSession> roster = sessions.roster();
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, clock.getAsLong()))
.toList();
return text(json(Map.of("leads", leadRows, "members", out)));
List<Map<String, Object>> capacity = roster.stream().map(MemberSession::profile).distinct().sorted()
.map(profile -> capacityView(profile, liveCount, maxLoad, roster, messages, clock.getAsLong()))
.toList();
return text(json(Map.of("leads", leadRows, "members", out, "capacity", capacity)));
} catch (HerdrException e) {
return error("herdr error listing the fleet: " + e.getMessage());
}
}
/**
* Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that bridged
* may stop the member: the bridge has capacity facts but no work list, and choosing work needs
* authority it does not have. {@code idleForSeconds} is derived from monotonic nanoTime and has
* no wall-clock meaning across a daemon restart.
*/
private static Map<String, Object> memberCapacityView(MemberSession session, Agent live,
MessageService messages, long nowNanos) {
Map<String, Object> row = SessionManager.rosterView(session, live);
boolean open = messages != null && messages.hasAcceptedDelivery(session.terminalId());
boolean inbox = messages != null && messages.hasInboxMessage(session.terminalId());
boolean reclaimable = (session.state() == MemberSession.State.READY || session.state() == MemberSession.State.DONE)
&& !open && !inbox;
row.put("reclaimable", reclaimable);
row.put("idleForSeconds", reclaimable ? Math.max(0, (nowNanos - session.lastActivityAtNanos()) / 1_000_000_000L) : null);
return row;
}
private static Map<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
Function<String, Integer> maxLoad, List<MemberSession> roster,
MessageService messages, long nowNanos) {
Integer cap = maxLoad.apply(profile);
int live = liveCount.apply(profile);
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
.filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE))
.filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId())))
.count();
Map<String, Object> row = new LinkedHashMap<>();
row.put("profile", profile); row.put("maxLoad", cap); row.put("live", live);
row.put("free", cap == null ? null : Math.max(0, cap - live)); row.put("reclaimable", reclaimable);
return row;
}
/**
* One lead's row: its address, its name, and whether it can be reached right now.
*
@@ -202,6 +202,16 @@ public final class MessageService {
return agents.status(target);
}
/** Read-only delegation fact for fleet views. */
public boolean hasAcceptedDelivery(String target) {
return rendezvous.isWaiting(target);
}
/** Read-only inbox fact for fleet views. */
public boolean hasInboxMessage(String target) {
return !inbox.peek(target).isEmpty();
}
/**
* Route a worker's explicit {@code bridge_reply}: resolve an open send, or queue it in the
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
@@ -342,6 +342,19 @@ class BridgeMcpTest {
assertTrue(out.contains("\"liveStatus\":\"unknown\""), out);
}
@Test
void capacityUsesThePlacementLiveCount() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
sessions.acquire("ltms-local", null, null, null);
McpSchema.CallToolResult res = BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, profile -> 2, profile -> 2, () -> 0, Map.of(), "");
String out = textOf(res);
assertTrue(out.contains("\"maxLoad\":2"), out);
assertTrue(out.contains("\"live\":2"), out);
assertTrue(out.contains("\"free\":0"), out);
}
@Test
void listReportsLeadsAndFlagsTheCallersOwnRow() {
FakeHerdr h = new FakeHerdr();
@@ -614,8 +627,8 @@ class BridgeMcpTest {
assertTrue(out.contains("\"role\":\"dev\""), out);
assertTrue(out.contains("\"role\":\"reviewer\""), out);
assertEquals(2, out.split("\"profile\":\"ltms-local\"", -1).length - 1,
"both members share one profile — that is the point: " + out);
assertEquals(3, out.split("\"profile\":\"ltms-local\"", -1).length - 1,
"two member rows and one capacity row share the profile: " + out);
}
@Test