Compare commits

..

5 Commits

Author SHA1 Message Date
Dai Ha 24559d81ac CB-573: require explicit capacity source
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 55s
2026-08-15 05:49:04 +02:00
Dai Ha 6c1c2c3994 CB-573: report empty configured profile capacity
CI / build (pull_request) Successful in 59s
CI / contract (pull_request) Successful in 1m0s
2026-08-15 05:45:39 +02:00
Dai Ha 01fab15713 CB-573: add fleet capacity view
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Successful in 1m25s
2026-08-15 05:43:46 +02:00
Dai Ha bf0ff2adbf CB-573: keep active delegations out of idle
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m27s
2026-08-15 05:40:37 +02:00
Dai Ha ed4bbc1c56 CB-573: add pure health classification model
CI / build (pull_request) Successful in 59s
CI / contract (pull_request) Successful in 1m15s
2026-08-15 05:37:30 +02:00
14 changed files with 321 additions and 253 deletions
@@ -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.
@@ -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));
}
}
@@ -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) { }
@@ -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);
}
@@ -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) { }
@@ -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
}
@@ -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<String, Integer> byTarget = new ConcurrentHashMap<>();
private final Map<String, Integer> 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); }
}
@@ -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<String, Long> lastProbe = new HashMap<>();
private int cursor;
public List<String> choose(List<String> candidates, long nowNanos, long configuredCooldownNanos) {
long cooldown = Math.max(COOLDOWN_NANOS, configuredCooldownNanos);
List<String> 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);
}
}
@@ -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<String, Integer> liveCount, Function<String, Integer> maxLoad,
Supplier<Set<String>> 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));
})
@@ -513,9 +515,6 @@ public final class BridgeMcp {
? "[done — worker finished without a structured bridge_reply; transcript tail follows]\n" + v.reply()
: v.reply());
case PENDING -> text("[pending — " + v.detail() + "]");
case ASKING -> text("[question — worker is waiting for your answer]\n" + v.reply()
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + v.turnId()
+ "\" and content set to your answer; the worker resumes the same turn.");
case FAILED -> text("[failed — " + v.detail() + "]");
};
}
@@ -722,6 +721,12 @@ public final class BridgeMcp {
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
Map<String, String> 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<String, String> leads, String selfTerm) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -731,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<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, capacity.clock().getAsLong()))
.toList();
return text(json(Map.of("leads", leadRows, "members", out)));
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
roster.stream().map(MemberSession::profile).forEach(profiles::add);
Map<String, Object> 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<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.
*
@@ -127,8 +127,6 @@ public final class MessageService {
public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */
PENDING,
/** The worker is paused in {@code bridge_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
ASKING,
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */
DONE,
/** The delegation could not complete (timed out, worker gone, or busy). */
@@ -138,28 +136,16 @@ public final class MessageService {
/**
* A poll snapshot of an async delegation.
*
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, or the question when
* {@link #phase} is {@link Phase#ASKING}; otherwise {@code null}
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, else {@code null}
* @param replySource {@code "reply"} (structured {@code bridge_reply}) or {@code "transcript"}
* (completion scrape) when {@link Phase#DONE}, else {@code null}
* @param detail a human note (live worker status while pending, ask state, or failure reason)
* @param turnId correlation id for an {@link Phase#ASKING} ticket, else {@code null}
* @param detail a human note (live worker status while pending, or the failure reason)
*/
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail,
String turnId) {
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail) {
}
/** An in-flight or finished async delegation, keyed by its ticket. */
private static final class Task {
private final String target;
private final CompletableFuture<Reply> future = new CompletableFuture<>();
private final long createdNanos = System.nanoTime();
private volatile Reply question;
private volatile String turnId;
private Task(String target) {
this.target = target;
}
private record Task(String target, CompletableFuture<Reply> future, long createdNanos) {
}
private final AgentControl agents;
@@ -170,10 +156,6 @@ public final class MessageService {
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
/** The async task that currently owns a target's send lock. */
private final ConcurrentHashMap<String, Task> asyncTasksByTarget = new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code bridge_ask} turn. */
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -220,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
@@ -351,9 +343,6 @@ public final class MessageService {
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
}
try {
if (hasAsyncQuestion(target)) {
return new Reply(Outcome.BUSY, null); // the worker's current turn is paused for its lead
}
// Open the waiter BEFORE queueing delivery (CB-548). A fast reply — the worker already
// injectable the instant we enqueue — otherwise arrives before the waiter is registered
// and orphans into the inbox while this send blocks to the timeout (the enqueue-before-
@@ -414,14 +403,12 @@ public final class MessageService {
rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
}
markAsyncQuestion(workerSession, question, ticket.turnId());
}
try {
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("bridge_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -464,12 +451,9 @@ public final class MessageService {
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
finishAsyncTask(turnId, result);
return result;
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
} catch (TimeoutException e) {
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
// turn (it never re-entered the injector), so a silent worker rides out the window.
@@ -509,27 +493,9 @@ public final class MessageService {
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(target);
tasks.put(ticket, task);
asyncExecutor.submit(() -> {
try {
Runnable trackingAccepted = () -> {
if (onAccepted != null) {
onAccepted.run();
}
asyncTasksByTarget.put(target, task);
};
Reply result = send(target, content, ASYNC_TIMEOUT_MS, trackingAccepted);
if (result.outcome() == Outcome.QUESTION) {
asyncTasksByTarget.remove(target, task);
} else {
finishAsyncTask(task, result);
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
asyncTasksByTarget.remove(target, task);
}
});
CompletableFuture<Reply> future = CompletableFuture.supplyAsync(
() -> send(target, content, ASYNC_TIMEOUT_MS, onAccepted), asyncExecutor);
tasks.put(ticket, new Task(target, future, System.nanoTime()));
pruneTerminalTickets();
log.debug("async send {} -> {}", ticket, target);
return ticket;
@@ -545,32 +511,27 @@ public final class MessageService {
if (task == null) {
return null;
}
CompletableFuture<Reply> f = task.future;
CompletableFuture<Reply> f = task.future();
if (!f.isDone()) {
Reply question = task.question;
if (question != null) {
return new TaskView(ticket, Phase.ASKING, question.text(), null,
"worker is waiting for your answer", question.turnId());
}
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target()));
}
Reply r;
try {
r = f.getNow(null);
} catch (CompletionException | java.util.concurrent.CancellationException e) {
Throwable cause = (e instanceof CompletionException ce && ce.getCause() != null) ? ce.getCause() : e;
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage(), null);
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage());
}
if (r.completed()) {
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null);
return new TaskView(ticket, Phase.DONE, r.text(), source, null);
}
// A wedged worker (CB-109) carries the error context as its reason; the timeout/busy
// outcomes carry none, so fall back to the outcome name.
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null
? r.text()
: "no reply — " + r.outcome().name().toLowerCase();
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
return new TaskView(ticket, Phase.FAILED, null, null, detail);
}
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
@@ -585,51 +546,7 @@ public final class MessageService {
/** Drop finished tickets older than the TTL so the registry cannot grow without bound. */
private void pruneTerminalTickets() {
long cutoff = System.nanoTime() - TICKET_TTL_NANOS;
tasks.values().removeIf(t -> t.future.isDone() && t.createdNanos < cutoff);
}
/** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */
private void markAsyncQuestion(String target, String text, String turnId) {
Task task = asyncTasksByTarget.get(target);
if (task != null) {
task.question = new Reply(Outcome.QUESTION, text, turnId);
task.turnId = turnId;
asyncTasksByTurn.put(turnId, task);
}
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.question = null;
if (forgetTurn) {
asyncTasksByTurn.remove(turnId, task);
task.turnId = null;
}
}
}
/** Complete and detach an async ticket after its worker's actual terminal reply. */
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
asyncTasksByTarget.remove(task.target, task);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
}
/** Complete the async ticket correlated to a specific answered turn. */
private void finishAsyncTask(String turnId, Reply result) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
finishAsyncTask(task, result);
}
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
tasks.values().removeIf(t -> t.future().isDone() && t.createdNanos() < cutoff);
}
/** Release the async executor. */
@@ -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);
}
}
@@ -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<String> 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));
}
}
@@ -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;
}
@@ -132,131 +132,6 @@ class BridgeMcpTest {
assertEquals("async LGTM", textOf(polled));
}
@Test
void asyncSendSurfacesAnAskThenKeepsTheTicketForTheFinalReply() throws Exception {
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
McpSchema.CallToolResult question = BridgeMcp.poll(messages, ticket, null);
deadline = System.currentTimeMillis() + 3000;
while (!textOf(question).contains("[question") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
question = BridgeMcp.poll(messages, ticket, null);
}
assertTrue(textOf(question).contains("which config?"), textOf(question));
String questionText = textOf(question);
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
BridgeMcp.reply(messages, "term_a", "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
McpSchema.CallToolResult done = BridgeMcp.poll(messages, ticket, null);
deadline = System.currentTimeMillis() + 3000;
while (!"done".equals(textOf(done)) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
done = BridgeMcp.poll(messages, ticket, null);
}
assertEquals("done", textOf(done));
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
McpSchema.CallToolResult ask = BridgeMcp.ask(messages, "term_a", "still there?", 50L);
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(BridgeMcp.poll(messages, ticket, null)).startsWith("[pending"));
BridgeMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
}
@Test
void asyncSendFailureDoesNotLeaveItsTicketPending() throws Exception {
String ticket = messages.sendAsync("term_a", "do it", () -> {
throw new IllegalStateException("accept failed");
});
long deadline = System.currentTimeMillis() + 3000;
MessageService.TaskView view = messages.poll(ticket);
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
view = messages.poll(ticket);
}
assertEquals(MessageService.Phase.FAILED, view.phase());
assertEquals("accept failed", view.detail());
}
@Test
void anotherAsyncTicketCannotCaptureAReplyWhileTheFirstTicketIsAsking() throws Exception {
McpSchema.CallToolResult firstAccepted = BridgeMcp.sendAsync(messages, "term_a", "first", null, Set.of());
String firstTicket = textOf(firstAccepted).substring(textOf(firstAccepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
MessageService.TaskView first = messages.poll(firstTicket);
deadline = System.currentTimeMillis() + 3000;
while (first.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
first = messages.poll(firstTicket);
}
assertEquals(MessageService.Phase.ASKING, first.phase());
String firstTurnId = first.turnId();
McpSchema.CallToolResult secondAccepted = BridgeMcp.sendAsync(messages, "term_a", "second", null, Set.of());
String secondTicket = textOf(secondAccepted).substring(textOf(secondAccepted).indexOf("ticket=") + "ticket=".length()).trim();
MessageService.TaskView second = messages.poll(secondTicket);
deadline = System.currentTimeMillis() + 3000;
while (second.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
second = messages.poll(secondTicket);
}
assertEquals(MessageService.Phase.FAILED, second.phase());
BridgeMcp.reply(messages, "term_a", "late reply");
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, firstTurnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
BridgeMcp.reply(messages, "term_a", "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@Test
void pollUnknownTicketIsAnError() {
McpSchema.CallToolResult res = BridgeMcp.poll(messages, "task-999", null);
@@ -467,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();
@@ -740,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