diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 7a16453..85d057c 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, new BridgeMcp.CapacitySource(profile -> liveCountRef.get().apply(profile), + profile -> { + var configured = config.get().profiles().get(profile); + return configured == null ? null : configured.maxLoad(); + }, () -> config.get().profiles().keySet(), 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/health/FleetHealth.java b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealth.java new file mode 100644 index 0000000..5c8ca23 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealth.java @@ -0,0 +1,40 @@ +package dev.ltms.bridged.health; + +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.session.MemberSession; + +/** + * Pure classifier. Collection and repair are deliberately outside this package. + * {@link HealthState#ERROR_ON_SCREEN} is not decided yet because it needs a bounded pane detection + * read and an adapter-specific fatal signature; status facts alone must not guess it. + */ +public final class FleetHealth { + private FleetHealth() { } + + public static HealthDecision decide(HealthSnapshot s, HealthPrior prior, long nowNanos) { + if (s.controlLinkDown()) return result(HealthState.CONTROL_LINK_DOWN, false); + if (s.targetNotFound()) return result(HealthState.GONE, false); + if (s.sessionState() == MemberSession.State.SPAWNING && !s.present() && s.readinessGraceElapsed()) { + return result(HealthState.NEVER_READY, false); + } + if (s.orphanedDelegation()) return result(HealthState.DELEGATION_ORPHANED, false); + boolean disagreement = s.sessionState() == MemberSession.State.BUSY && s.acceptedDelivery() + && (s.liveStatus() == AgentStatus.IDLE || s.liveStatus() == AgentStatus.DONE); + if (disagreement && prior.busyButDone()) return result(HealthState.TURN_BOUNDARY_LOST, true); + if (s.stalled()) return result(HealthState.STALL_SUSPECTED, disagreement); + if (s.replyStranded()) return result(HealthState.REPLY_STRANDED, disagreement); + if (s.queuedDelivery() || s.inboxMessage()) return result(HealthState.WORK_PENDING, disagreement); + if (s.sessionState() == MemberSession.State.SPAWNING) return result(HealthState.STARTING, disagreement); + if (s.acceptedDelivery() && s.liveStatus() == AgentStatus.BLOCKED) { + return result(HealthState.BLOCKED_AMBIGUOUS, disagreement); + } + // An accepted delivery remains bridge work even when herdr is late, unknown, or has already + // reported DONE once. It cannot be IDLE until the delegation has resolved. + if (s.acceptedDelivery()) return result(HealthState.WORKING, disagreement); + return result(HealthState.IDLE, disagreement); + } + + private static HealthDecision result(HealthState state, boolean disagreement) { + return new HealthDecision(state, new HealthPrior(disagreement)); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/health/HealthDecision.java b/bridged/src/main/java/dev/ltms/bridged/health/HealthDecision.java new file mode 100644 index 0000000..fc10ed0 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/HealthDecision.java @@ -0,0 +1,4 @@ +package dev.ltms.bridged.health; + +/** Classification plus the private fact that the next pure decision needs. */ +public record HealthDecision(HealthState state, HealthPrior prior) { } diff --git a/bridged/src/main/java/dev/ltms/bridged/health/HealthPrior.java b/bridged/src/main/java/dev/ltms/bridged/health/HealthPrior.java new file mode 100644 index 0000000..6c0816e --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/HealthPrior.java @@ -0,0 +1,6 @@ +package dev.ltms.bridged.health; + +/** Private cross-tick observation. It is deliberately not a reported health value. */ +public record HealthPrior(boolean busyButDone) { + public static final HealthPrior NONE = new HealthPrior(false); +} diff --git a/bridged/src/main/java/dev/ltms/bridged/health/HealthSnapshot.java b/bridged/src/main/java/dev/ltms/bridged/health/HealthSnapshot.java new file mode 100644 index 0000000..0089c28 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/HealthSnapshot.java @@ -0,0 +1,11 @@ +package dev.ltms.bridged.health; + +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.session.MemberSession; + +/** Read-only facts from one fleet collection tick. */ +public record HealthSnapshot(MemberSession.State sessionState, AgentStatus liveStatus, + boolean acceptedDelivery, boolean queuedDelivery, boolean inboxMessage, + boolean present, boolean targetNotFound, boolean controlLinkDown, + boolean readinessGraceElapsed, boolean orphanedDelegation, + boolean replyStranded, boolean stalled) { } diff --git a/bridged/src/main/java/dev/ltms/bridged/health/HealthState.java b/bridged/src/main/java/dev/ltms/bridged/health/HealthState.java new file mode 100644 index 0000000..5a9c0bd --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/HealthState.java @@ -0,0 +1,8 @@ +package dev.ltms.bridged.health; + +/** Health classifications reported for a member. */ +public enum HealthState { + STARTING, IDLE, WORKING, WORK_PENDING, BLOCKED_AMBIGUOUS, + NEVER_READY, GONE, TURN_BOUNDARY_LOST, ERROR_ON_SCREEN, STALL_SUSPECTED, + MUTE, REPLY_STRANDED, DELEGATION_ORPHANED, CONTROL_LINK_DOWN +} diff --git a/bridged/src/main/java/dev/ltms/bridged/health/MuteCounter.java b/bridged/src/main/java/dev/ltms/bridged/health/MuteCounter.java new file mode 100644 index 0000000..72f2c76 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/MuteCounter.java @@ -0,0 +1,25 @@ +package dev.ltms.bridged.health; + +import dev.ltms.bridged.msg.Rendezvous; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +/** + * Counts turns that ended via the completion fallback instead of {@code bridge_reply}. + * MUTE is an observation by target and profile, not a classifier state and never suppresses faults. + */ +public final class MuteCounter { + private final Map byTarget = new ConcurrentHashMap<>(); + private final Map byProfile = new ConcurrentHashMap<>(); + + /** Record only fallback completion; a structured reply does not make a member mute. */ + public void observe(String target, String profile, Rendezvous.Kind kind) { + if (kind != Rendezvous.Kind.COMPLETION) return; + byTarget.merge(target, 1, Integer::sum); + byProfile.merge(profile, 1, Integer::sum); + } + + public int forTarget(String target) { return byTarget.getOrDefault(target, 0); } + public int forProfile(String profile) { return byProfile.getOrDefault(profile, 0); } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/health/PaneBudget.java b/bridged/src/main/java/dev/ltms/bridged/health/PaneBudget.java new file mode 100644 index 0000000..315c76a --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/health/PaneBudget.java @@ -0,0 +1,26 @@ +package dev.ltms.bridged.health; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +/** Fixed pane-probe limits. Pane content is never retained here. */ +public final class PaneBudget { + public static final long COOLDOWN_NANOS = 60_000_000_000L; + public static final int MAX_PER_TICK = 2; + private final Map lastProbe = new HashMap<>(); + private int cursor; + + public List choose(List candidates, long nowNanos, long configuredCooldownNanos) { + long cooldown = Math.max(COOLDOWN_NANOS, configuredCooldownNanos); + List out = new ArrayList<>(); + for (int n = 0; n < candidates.size() && out.size() < MAX_PER_TICK; n++) { + String target = candidates.get((cursor + n) % candidates.size()); + Long last = lastProbe.get(target); + if (last == null || nowNanos - last >= cooldown) { out.add(target); lastProbe.put(target, nowNanos); } + } + if (!candidates.isEmpty()) cursor = (cursor + 1) % candidates.size(); + return List.copyOf(out); + } +} 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..0c960e4 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,8 @@ 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.function.Supplier; import java.util.stream.Collectors; /** @@ -74,15 +76,14 @@ 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 CapacitySource capacity; - /** - * Legacy constructor — no authorization. Retained so existing tests exercise tool behaviour - * without an auth fixture. - */ - public BridgeMcp(MessageService messages, PeerLauncher workers, - SessionManager sessions, ConnectionIdentity identity, MemberPresence presence, - PrimaryRegistry primaryRegistry) { - this(messages, workers, sessions, identity, presence, primaryRegistry, null, null); + /** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */ + public record CapacitySource(Function liveCount, Function maxLoad, + Supplier> configuredProfiles, LongSupplier clock) { + /** Inert test-only source. It omits capacity rather than inventing zero live counts. */ + public static CapacitySource none() { return new CapacitySource(_ -> 0, _ -> null, Set::of, System::nanoTime); } + boolean available() { return !configuredProfiles.get().isEmpty(); } } /** @@ -92,9 +93,10 @@ public final class BridgeMcp { * filter, so the REST guard does not cover it. * @param metrics registry for auth-failure counting; may be {@code null} */ - public BridgeMcp(MessageService messages, PeerLauncher workers, - SessionManager sessions, ConnectionIdentity identity, MemberPresence presence, - PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics) { + public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, Metrics metrics, CapacitySource capacity) { + this.capacity = capacity; McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() .jsonMapper(json) @@ -208,7 +210,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, capacity, callers == null ? Map.of() : callers.leads(), callerTerminal(exchange)); }) @@ -719,6 +721,12 @@ public final class BridgeMcp { */ static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, Map leads, String selfTerm) { + return listFleet(workers, sessions, null, CapacitySource.none(), leads, selfTerm); + } + + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, + Map leads, String selfTerm) { try { Map live = workers.list().stream() .map(Agent.class::cast) @@ -728,15 +736,56 @@ 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, capacity.clock().getAsLong())) .toList(); - return text(json(Map.of("leads", leadRows, "members", out))); + Set profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get()); + roster.stream().map(MemberSession::profile).forEach(profiles::add); + Map result = new LinkedHashMap<>(); + result.put("leads", leadRows); result.put("members", out); + if (capacity.available()) result.put("capacity", profiles.stream() + .map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages, + capacity.clock().getAsLong())).toList()); + return text(json(result)); } 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/health/FleetHealthTest.java b/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthTest.java new file mode 100644 index 0000000..6cc2e55 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthTest.java @@ -0,0 +1,56 @@ +package dev.ltms.bridged.health; + +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.session.MemberSession; +import dev.ltms.bridged.msg.Rendezvous; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class FleetHealthTest { + @Test void muteCountsOnlyCompletionFallbacks() { + MuteCounter mute = new MuteCounter(); + mute.observe("target", "terra", Rendezvous.Kind.REPLY); + mute.observe("target", "terra", Rendezvous.Kind.COMPLETION); + assertEquals(1, mute.forTarget("target")); + assertEquals(1, mute.forProfile("terra")); + } + @Test void turnBoundaryNeedsTwoSnapshots() { + HealthSnapshot s = snapshot(MemberSession.State.BUSY, AgentStatus.DONE, true); + HealthDecision first = FleetHealth.decide(s, HealthPrior.NONE, 1); + assertEquals(HealthState.WORKING, first.state()); + assertEquals(new HealthPrior(true), first.prior()); + assertEquals(HealthState.TURN_BOUNDARY_LOST, FleetHealth.decide(s, first.prior(), 2).state()); + } + + @Test void unknownLiveStatusWithAcceptedDeliveryIsNotIdle() { + assertEquals(HealthState.WORKING, FleetHealth.decide( + snapshot(MemberSession.State.BUSY, AgentStatus.UNKNOWN, true), HealthPrior.NONE, 1).state()); + } + + @Test void acceptedDeliveryNeverReportsIdle() { + for (MemberSession.State session : MemberSession.State.values()) { + for (AgentStatus live : AgentStatus.values()) { + HealthSnapshot s = snapshot(session, live, true); + assertEquals(false, FleetHealth.decide(s, HealthPrior.NONE, 1).state() == HealthState.IDLE, + () -> "accepted delivery returned IDLE for " + session + "/" + live); + } + } + } + + @Test void blockedDoesNotGuessPromptKind() { + assertEquals(HealthState.BLOCKED_AMBIGUOUS, FleetHealth.decide( + snapshot(MemberSession.State.BUSY, AgentStatus.BLOCKED, true), HealthPrior.NONE, 1).state()); + } + + @Test void controlLinkOutranksMemberFault() { + HealthSnapshot s = new HealthSnapshot(MemberSession.State.BUSY, AgentStatus.DONE, true, false, + false, true, true, true, false, true, true, true); + assertEquals(HealthState.CONTROL_LINK_DOWN, FleetHealth.decide(s, HealthPrior.NONE, 1).state()); + } + + private static HealthSnapshot snapshot(MemberSession.State state, AgentStatus live, boolean accepted) { + return new HealthSnapshot(state, live, accepted, false, false, true, false, false, + false, false, false, false); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/health/PaneBudgetTest.java b/bridged/src/test/java/dev/ltms/bridged/health/PaneBudgetTest.java new file mode 100644 index 0000000..162ec20 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/health/PaneBudgetTest.java @@ -0,0 +1,14 @@ +package dev.ltms.bridged.health; + +import org.junit.jupiter.api.Test; +import java.util.List; +import static org.junit.jupiter.api.Assertions.assertEquals; + +class PaneBudgetTest { + @Test void fixedLimitsIgnoreWeakerConfig() { + PaneBudget budget = new PaneBudget(); + List targets = List.of("a", "b", "c"); + assertEquals(List.of("a", "b"), budget.choose(targets, 0, 0)); + assertEquals(List.of("c"), budget.choose(targets, 1, 0)); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java index 96923cf..fb2e067 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpAuthzTest.java @@ -72,7 +72,7 @@ class BridgeMcpAuthzTest { new PrimaryRegistry(null), enforce ? CallerResolver.withLeadsAndMembers(identity, false, null, Map::of, new MemberRegistry(null)) : null, - metrics); + metrics, BridgeMcp.CapacitySource.none()); return mcp; } 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..b5b6b00 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,42 @@ 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, new BridgeMcp.CapacitySource(profile -> 2, profile -> 2, + () -> Set.of("ltms-local"), () -> 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 capacityIncludesConfiguredProfileWithoutMembers() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), + sessions, null, new BridgeMcp.CapacitySource(profile -> 0, profile -> 2, + () -> Set.of("terra"), () -> 0), Map.of(), "")); + assertTrue(out.contains("\"profile\":\"terra\""), out); + assertTrue(out.contains("\"live\":0"), out); + assertTrue(out.contains("\"free\":2"), out); + assertTrue(out.contains("\"reclaimable\":0"), out); + } + + @Test + void inertCapacitySourceOmitsCapacityBlock() { + FakeHerdr h = new FakeHerdr(); + String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), + new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))), null, + BridgeMcp.CapacitySource.none(), Map.of(), "")); + assertFalse(out.contains("\"capacity\":"), out); + } + @Test void listReportsLeadsAndFlagsTheCallersOwnRow() { FakeHerdr h = new FakeHerdr(); @@ -615,7 +651,7 @@ 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); + "inert capacity is omitted, leaving the two member rows: " + out); } @Test