Compare commits

..

9 Commits

Author SHA1 Message Date
Dai Ha a3463264f8 CB-568: fail queued async tickets on teardown
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m26s
2026-08-15 06:10:24 +02:00
Dai Ha 16e17b32ad CB-568: preserve dropped turn causes
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 1m26s
2026-08-15 06:06:07 +02:00
Dai Ha 8a53d5bfc6 Merge CB-574: an async delegation can now receive a worker's question
CI / contract (push) Successful in 1m2s
CI / build (push) Failing after 1m40s
A worker on a wait:false delegation called bridge_ask and the lead never
saw the question. Outcome.QUESTION is deliberately non-terminal, but
taskView tested r.completed() and fell into the failure branch, so the
ticket was marked FAILED and both the question text and its turnId were
discarded. The worker blocked for 55s, gave up, and had to abandon its
task. CLAUDE.md tells leads to prefer wait:false and to answer an ask with
bridge_send{turnId, content}; those two could not both be followed.

bridge_poll now returns a non-terminal ASKING phase carrying the question
and its turnId, and the ticket stays live so the worker's real reply still
lands on it. An unanswered ask returns the ticket to PENDING, because only
the question wait ended - the delegated turn continues. The 55s/115s ask
caps are unchanged: they exist because the worker's own MCP call would time
out, so widening them would only move the failure.

Two defects found reviewing the first revision, both from replacing
supplyAsync with a manually completed future:

- an exception inside the send left the future uncompleted, so the ticket
  stayed PENDING for the life of the daemon. Now caught and completed
  exceptionally.
- correlation was keyed by target, one entry per worker, registered before
  the session lock. With two tickets outstanding on one target the second
  overwrote the first, so a late reply could resolve the wrong ticket.
  Correlation is now per turn, the target entry exists only while that send
  owns the lock, and a reply with no live waiter still goes to the durable
  inbox as before.
2026-08-15 05:50:45 +02:00
Dai Ha 695da7418e Merge CB-573 (part 1): health classification model and the bridge_list capacity view
CI / contract (push) Successful in 45s
CI / build (push) Successful in 55s
Fleet capacity was invisible. A finished member held a terra slot until a
spawn was refused with 'at maxLoad: 2 live >= 2 cap', and nothing had told
the lead the slot was still held. bridge_list now reports, per configured
profile, maxLoad / live / free / reclaimable, and per member idleForSeconds
and reclaimable.

live comes from the same liveCountRef function placement consumes, so the
advertised free slots cannot drift from what bridge_spawn will accept. The
profile list is the union of configured and roster profiles: an empty
configured profile still appears with its full capacity, and a member whose
profile was removed from config stays visible rather than vanishing.

reclaimable is advisory. The bridge never spawns, stops or retasks a member
to improve utilisation: it has capacity facts but no work list, and choosing
work needs authority it does not have.

Also lands the pure health classifier, its precedence chain, the MUTE counter
and the pane budget. The classifier never reports IDLE while an accepted
delivery is open — IDLE is a claim that nothing is outstanding, and the
capacity view reads exactly that field.

Capacity dependencies are one required CapacitySource rather than defaulted
constructor arguments. A defaulted liveCount would report free slots that do
not exist, which is the dangerous direction; CapacitySource.none() omits the
block instead of inventing zeros.
2026-08-15 05:49:51 +02:00
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
19 changed files with 493 additions and 45 deletions
@@ -285,6 +285,12 @@ public final class Bridged {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
}
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
presence::forget);
@@ -383,7 +389,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);
}
}
@@ -137,7 +137,13 @@ public final class CompletionResolver implements TurnListener {
@Override
public void onTurnFailed(String target) {
InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn));
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn, null));
}
@Override
public void onTurnFailed(String target, String reason) {
InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn, reason));
}
/** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */
@@ -192,6 +198,11 @@ public final class CompletionResolver implements TurnListener {
/** Synchronous fail (the unit-testable core of {@link #onTurnFailed}). */
void fail(String target, InFlight turn) {
fail(target, turn, null);
}
/** Synchronous fail with an optional reason supplied by a dropped worker queue. */
void fail(String target, InFlight turn, String explicitReason) {
// A never-delivered readiness failure has no in-flight record but still has a blocked send;
// fall back to the currently-registered waiter (unambiguous — that send never completed, so
// no next turn exists to confuse it with).
@@ -201,16 +212,18 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn); // nobody blocked on this worker — nothing to fail
return;
}
String reason;
try {
reason = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
reason = "";
}
if (reason.isBlank()) {
// No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110).
reason = "worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)";
String reason = explicitReason;
if (reason == null || reason.isBlank()) {
try {
reason = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
reason = "";
}
if (reason.isBlank()) {
// No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110).
reason = "worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)";
}
}
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
@@ -410,8 +410,8 @@ public final class Injector {
for (Pending p : pending) {
p.delivered().completeExceptionally(cause);
}
if (hadDeliveredTurn) {
turnListener.onTurnFailed(target);
}
// A queued send has no in-flight record, while a delivered turn does. CompletionResolver
// handles both forms and resolves its waiter at most once.
turnListener.onTurnFailed(target, cause.getMessage());
}
}
@@ -41,6 +41,14 @@ public interface TurnListener {
default void onTurnFailed(String target) {
}
/**
* As {@link #onTurnFailed(String)}, carrying the reason a worker became unreachable. The default
* keeps existing listeners working while allowing the completion resolver to report a useful cause.
*/
default void onTurnFailed(String target, String reason) {
onTurnFailed(target);
}
/**
* A message was just delivered into {@code target}'s pane (CB-115). Fired so the completion
* resolver can snapshot the pane's pre-turn content: a later {@link #onTurnComplete} whose
@@ -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));
})
@@ -722,6 +724,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 +739,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.
*
@@ -9,6 +9,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
@@ -170,8 +171,8 @@ 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 tasks that have accepted delivery for a target. */
private final ConcurrentHashMap<String, Set<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();
@@ -220,6 +221,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
@@ -290,14 +301,14 @@ public final class MessageService {
*/
public boolean abandon(String target, String reason) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
if (waiter == null || waiter.isDone()) {
return false; // nobody is blocked on this worker — nothing to abandon
}
boolean failed = rendezvous.resolveFailure(waiter, reason);
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
boolean asyncFailed = tasks.values().stream()
.filter(task -> target.equals(task.target) && task.question == null)
.anyMatch(task -> task.future.complete(new Reply(Outcome.WORKER_FAILED, reason)));
if (failed) {
log.warn("abandoning the blocked send to {}: {}", target, reason);
}
return failed;
return failed || asyncFailed;
}
/**
@@ -344,6 +355,11 @@ public final class MessageService {
* never earned. {@code null} disables the hook.
*/
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
return send(target, content, timeoutMillis, onAccepted, null);
}
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
@@ -351,6 +367,9 @@ public final class MessageService {
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
}
try {
if (task != null && task.future.isDone()) {
return task.future.getNow(null);
}
if (hasAsyncQuestion(target)) {
return new Reply(Outcome.BUSY, null); // the worker's current turn is paused for its lead
}
@@ -517,17 +536,18 @@ public final class MessageService {
if (onAccepted != null) {
onAccepted.run();
}
asyncTasksByTarget.put(target, task);
asyncTasksByTarget.computeIfAbsent(target, _ -> ConcurrentHashMap.newKeySet()).add(task);
};
Reply result = send(target, content, ASYNC_TIMEOUT_MS, trackingAccepted);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, trackingAccepted, task);
if (result.outcome() == Outcome.QUESTION) {
asyncTasksByTarget.remove(target, task);
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
} else {
finishAsyncTask(task, result);
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
asyncTasksByTarget.remove(target, task);
untrackAsyncTarget(task);
}
});
pruneTerminalTickets();
@@ -590,7 +610,11 @@ public final class MessageService {
/** 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);
Set<Task> targetTasks = asyncTasksByTarget.get(target);
Task task = targetTasks == null ? null : targetTasks.stream()
.filter(candidate -> !candidate.future.isDone())
.findFirst()
.orElse(null);
if (task != null) {
task.question = new Reply(Outcome.QUESTION, text, turnId);
task.turnId = turnId;
@@ -613,7 +637,7 @@ public final class MessageService {
/** 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);
untrackAsyncTarget(task);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
@@ -632,6 +656,17 @@ public final class MessageService {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
}
/** Stop tracking a task once it no longer owns an accepted target turn. */
private void untrackAsyncTarget(Task task) {
Set<Task> targetTasks = asyncTasksByTarget.get(task.target);
if (targetTasks != null) {
targetTasks.remove(task);
if (targetTasks.isEmpty()) {
asyncTasksByTarget.remove(task.target, targetTasks);
}
}
}
/** Release the async executor. */
public void close() {
asyncExecutor.shutdown();
@@ -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));
}
}
@@ -259,6 +259,7 @@ class InjectorTest {
private static final class Captor implements TurnListener {
final List<String> completed = new ArrayList<>();
final List<String> failed = new ArrayList<>();
final List<String> failureReasons = new ArrayList<>();
@Override
public void onTurnComplete(String target) {
@@ -269,6 +270,12 @@ class InjectorTest {
public void onTurnFailed(String target) {
failed.add(target);
}
@Override
public void onTurnFailed(String target, String reason) {
failed.add(target);
failureReasons.add(reason);
}
}
// ~30s of unknown at the 250ms prod poll interval; enough onStatus samples to trip the stall.
@@ -322,6 +329,23 @@ class InjectorTest {
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
}
@Test
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered");
CompletableFuture<Void> queued = inj.enqueue(T, "queued");
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
inj.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
assertEquals(List.of(T), cap.failed, "drop signals one turn failure for both affected states");
assertEquals(List.of("agent target sol not found"), cap.failureReasons);
assertTrue(delivered.isDone(), "the delivered future has already completed");
assertTrue(queued.isCompletedExceptionally(), "the queued future fails with the drop cause");
}
@Test
void dropFailsTheTurnOfADeliveredMessageWhenTheWorkerVanishes() {
// CB-110: the message was delivered (no longer queued), so failing queued waiters alone would
@@ -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;
}
@@ -467,6 +467,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 +776,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
@@ -120,6 +120,36 @@ class MessageServiceTest {
assertFalse(reply.completed());
}
@Test
void droppedQueuedAndDeliveredTurnsExposeTheRealCauseExactlyOnce() throws Exception {
CompletableFuture<MessageService.Reply> first = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // first delivery
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
CompletableFuture<Void> queued = injector.enqueue(T, "second task");
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
MessageService.Reply reply = first.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome());
assertEquals("agent target sol not found", reply.text());
assertTrue(queued.isCompletedExceptionally(), "the queued delivery future also fails");
assertFalse(rendezvous.resolveFailure(waiter, "second failure"), "the waiter fails exactly once");
}
@Test
void noDropReasonKeepsTheExistingFallbackText() throws Exception {
herdr.readText("");
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
completion.onTurnFailed(T);
Rendezvous.Resolution resolution = waiter.get(2, TimeUnit.SECONDS);
assertEquals("worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)", resolution.text());
}
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
@Test
@@ -609,4 +639,57 @@ class MessageServiceTest {
assertTrue(view.detail() != null && view.detail().contains("released"),
"and the detail says why, rather than 'worker unknown'");
}
@Test
void abandonFailsEveryPendingAsyncTicketForTheReleasedTarget() throws Exception {
String first = messages.sendAsync(T, "first task");
awaitWaiting(); // first task owns the target lock and rendezvous waiter
String second = messages.sendAsync(T, "second task"); // parked on the same lock, not yet queued
assertTrue(messages.abandon(T, "agent target term_a not found"));
assertFailedTicket(first, "agent target term_a not found");
assertFailedTicket(second, "agent target term_a not found");
}
@Test
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertFalse(messages.abandon(T, "agent target term_a not found"),
"an asking ticket is an active turn, not a pending send to sweep");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
private void assertFailedTicket(String ticket, String reason) throws Exception {
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.FAILED);
assertEquals(reason, view.detail());
}
private MessageService.TaskView awaitTicketPhase(String ticket, MessageService.Phase phase) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
MessageService.TaskView view;
do {
view = messages.poll(ticket);
if (view.phase() == phase) {
return view;
}
Thread.sleep(5);
} while (System.currentTimeMillis() < deadline);
assertEquals(phase, view.phase());
return view;
}
}