diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 7a16453..8349fcb 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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. diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index dde9242..6c5a960 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -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 liveCount; + private final Function 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 liveCount, + Function 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 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 liveCount, Function maxLoad, + LongSupplier clock, Map leads, String selfTerm) { try { Map 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> out = sessions.roster().stream() - .map(s -> SessionManager.rosterView(s, live.get(s.terminalId()))) + List roster = sessions.roster(); + List> 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> 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 memberCapacityView(MemberSession session, Agent live, + MessageService messages, long nowNanos) { + Map 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 capacityView(String profile, Function liveCount, + Function maxLoad, List 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 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. * diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index e7758d4..0d0ab44 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -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 diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index 0aed9d2..7110233 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -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