Merge CB-580: fail a ticket when its member reaches a terminal health state
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m32s

Verified by the lead on 0af902e: mvn -f bridged/pom.xml clean install, unpiped, in a scratch
worktree — MVN_EXIT=0, Tests run: 689, Failures: 0, Errors: 0, BUILD SUCCESS.

Reviewed by the lead reading the diff. The member's own report was lost to the idle reaper before
it was collected, so there was no author write-up to review against.

All three defects that sank 3b2f395 are absent:
  * failTarget is required — one constructor, Objects.requireNonNull, no defaulting overload
    anywhere in the repo, and Bridged.java:365 updated to pass messages::abandon.
  * The failure fires inside reportTransition after its `if (previous == next) return;` guard, so an
    unchanged tick cannot reach it.
  * No AtomicReference; MessageService is passed directly, so there is no empty window.

abandon(String target, String reason) confirmed as CB-568's target-wide operation that resolves
waiters as a failure rather than letting them time out.

Known limitation, accepted and covered by its own test: states.put records the new state before the
bounded retries run, so if all three attempts throw, the tickets stay pending and no later tick
retries. It is logged at WARN, and the retries carry no backoff.
This commit was merged in pull request #56.
This commit is contained in:
2026-08-15 08:57:03 +02:00
3 changed files with 140 additions and 5 deletions
@@ -360,8 +360,10 @@ public final class Bridged {
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r -> var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-health-").unstarted(r)); Thread.ofVirtual().name("bridge-health-").unstarted(r));
if (cfg.health() != null && cfg.health().isEnabled()) { if (cfg.health() != null && cfg.health().isEnabled()) {
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
// through the same idempotent target-wide operation CB-516 already uses on release.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler, healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault()); System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true, String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured()); cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) { if ("detection-only".equals(coverage)) {
@@ -12,34 +12,50 @@ import java.util.HashMap;
import java.util.HashSet; import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier; import java.util.function.LongSupplier;
import java.util.function.Supplier; import java.util.function.Supplier;
/** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */ /** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */
public final class FleetHealthMonitor { public final class FleetHealthMonitor {
private static final Logger log = LoggerFactory.getLogger(FleetHealthMonitor.class); private static final Logger log = LoggerFactory.getLogger(FleetHealthMonitor.class);
/** Bounded attempts to run {@link #failTarget} for one transition. Never retried tick-to-tick (CB-580). */
static final int MAX_FAIL_TARGET_ATTEMPTS = 3;
private final AgentControl agents; private final AgentControl agents;
private final Supplier<List<MemberSession>> roster; private final Supplier<List<MemberSession>> roster;
private final MessageService messages; private final MessageService messages;
private final ScheduledExecutorService scheduler; private final ScheduledExecutorService scheduler;
private final LongSupplier clock; private final LongSupplier clock;
private final long intervalSeconds; private final long intervalSeconds;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>(); private final Map<String, HealthPrior> priors = new HashMap<>();
private final Map<String, HealthState> states = new HashMap<>(); private final Map<String, HealthState> states = new HashMap<>();
// These facts need the evidence publishers introduced by later M4 units. They are not negatives. // These facts need the evidence publishers introduced by later M4 units. They are not negatives.
private static final boolean NOT_YET_OBSERVED = false; private static final boolean NOT_YET_OBSERVED = false;
/**
* @param failTarget CB-568's idempotent target-wide failure operation (e.g. {@code messages::abandon}),
* invoked once when a member transitions into a terminal health state. Required —
* there is deliberately no defaulting overload; a caller that does not want the
* fail-tickets-on-terminal-health behavior must pass an explicit inert value (see
* {@code TestTurnTokens.inert} / {@code BridgeMcp.CapacitySource.none()} for the pattern).
*/
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages, public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) { ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
BiConsumer<String, String> failTarget) {
this.agents = agents; this.agents = agents;
this.roster = roster; this.roster = roster;
this.messages = messages; this.messages = messages;
this.scheduler = scheduler; this.scheduler = scheduler;
this.clock = clock; this.clock = clock;
this.intervalSeconds = intervalSeconds; this.intervalSeconds = intervalSeconds;
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
} }
/** Pure per-member decision seam. */ /** Pure per-member decision seam. */
@@ -91,6 +107,34 @@ public final class FleetHealthMonitor {
} else if (previous != null && fault(previous)) { } else if (previous != null && fault(previous)) {
log.info("fleet health member={} recovered state={} previous={}", target, next, previous); log.info("fleet health member={} recovered state={} previous={}", target, next, previous);
} }
// CB-580: a member entering GONE/NEVER_READY must not leave its waiting tickets pending
// forever. Fire exactly once per transition — never on a tick where the state is unchanged,
// which is what made the rejected commit call abandon() once per tick for as long as a
// member stayed terminal.
if (terminal(next)) {
failTerminalTarget(target, next);
}
}
private void failTerminalTarget(String target, HealthState state) {
String reason = "fleet health: member reached terminal state " + state.name();
RuntimeException last = null;
for (int attempt = 1; attempt <= MAX_FAIL_TARGET_ATTEMPTS; attempt++) {
try {
failTarget.accept(target, reason);
return;
} catch (RuntimeException error) {
last = error;
log.warn("fleet health: failTarget attempt {}/{} failed for member={} state={}",
attempt, MAX_FAIL_TARGET_ATTEMPTS, target, state, error);
}
}
log.warn("fleet health: giving up on failTarget for member={} state={} after {} attempts",
target, state, MAX_FAIL_TARGET_ATTEMPTS, last);
}
private static boolean terminal(HealthState state) {
return state == HealthState.GONE || state == HealthState.NEVER_READY;
} }
private static boolean fault(HealthState state) { private static boolean fault(HealthState state) {
@@ -20,8 +20,10 @@ import org.slf4j.LoggerFactory;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.function.BiConsumer;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class FleetHealthMonitorTest { class FleetHealthMonitorTest {
@Test void oneTickUsesOneFleetListForAnyRosterSize() { @Test void oneTickUsesOneFleetListForAnyRosterSize() {
@@ -37,7 +39,8 @@ class FleetHealthMonitorTest {
herdr.calls.clear(); herdr.calls.clear();
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()); MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
var scheduler = Executors.newSingleThreadScheduledExecutor(); var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60); FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60,
(_, _) -> { });
monitor.tick(); monitor.tick();
monitor.stop(); monitor.stop();
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count()); assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
@@ -49,7 +52,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor(); var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of, FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()), new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60); scheduler, () -> 1, 60, (_, _) -> { });
monitor.tick(); monitor.tick();
herdr.healthy(true); herdr.healthy(true);
monitor.tick(); monitor.tick();
@@ -68,7 +71,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor(); var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of, FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()), new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60); scheduler, () -> 1, 60, (_, _) -> { });
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST); monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST); monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.stop(); monitor.stop();
@@ -78,4 +81,90 @@ class FleetHealthMonitorTest {
logger.detachAppender(appender); logger.detachAppender(appender);
} }
} }
// --- CB-580: a member that reaches GONE/NEVER_READY must fail its waiting tickets
private static FleetHealthMonitor monitorWith(BiConsumer<String, String> failTarget) {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
var scheduler = Executors.newSingleThreadScheduledExecutor();
return new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, failTarget);
}
@Test void terminalTransitionFailsTheTargetOnce() {
RecordingFailTarget failTarget = new RecordingFailTarget();
FleetHealthMonitor monitor = monitorWith(failTarget);
monitor.reportTransition("term_a", HealthState.GONE);
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertEquals("term_a", failTarget.calls.get(0).target());
assertTrue(failTarget.calls.get(0).reason().contains("GONE"));
}
@Test void neverReadyNamesItselfAsTheReason() {
RecordingFailTarget failTarget = new RecordingFailTarget();
FleetHealthMonitor monitor = monitorWith(failTarget);
monitor.reportTransition("term_a", HealthState.NEVER_READY);
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertTrue(failTarget.calls.get(0).reason().contains("NEVER_READY"));
}
@Test void stayingInATerminalStateProducesOneFailureNotN() {
RecordingFailTarget failTarget = new RecordingFailTarget();
FleetHealthMonitor monitor = monitorWith(failTarget);
monitor.reportTransition("term_a", HealthState.GONE);
monitor.reportTransition("term_a", HealthState.GONE);
monitor.reportTransition("term_a", HealthState.GONE);
monitor.reportTransition("term_a", HealthState.GONE);
monitor.stop();
assertEquals(1, failTarget.calls.size());
}
@Test void aNonTerminalFaultStateDoesNotFailTheTarget() {
RecordingFailTarget failTarget = new RecordingFailTarget();
FleetHealthMonitor monitor = monitorWith(failTarget);
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.stop();
assertEquals(0, failTarget.calls.size());
}
@Test void failTargetRetryIsBounded() {
AlwaysThrowingFailTarget failTarget = new AlwaysThrowingFailTarget();
FleetHealthMonitor monitor = monitorWith(failTarget);
monitor.reportTransition("term_a", HealthState.GONE);
monitor.stop();
assertEquals(FleetHealthMonitor.MAX_FAIL_TARGET_ATTEMPTS, failTarget.calls);
}
@Test void exhaustedRetryStillDoesNotRefireOnAnUnchangedTick() {
AlwaysThrowingFailTarget failTarget = new AlwaysThrowingFailTarget();
FleetHealthMonitor monitor = monitorWith(failTarget);
monitor.reportTransition("term_a", HealthState.GONE);
int afterFirstTransition = failTarget.calls;
monitor.reportTransition("term_a", HealthState.GONE);
monitor.stop();
assertEquals(afterFirstTransition, failTarget.calls);
}
private record RecordedCall(String target, String reason) { }
private static final class RecordingFailTarget implements BiConsumer<String, String> {
final java.util.List<RecordedCall> calls = new java.util.ArrayList<>();
@Override public void accept(String target, String reason) {
calls.add(new RecordedCall(target, reason));
}
}
private static final class AlwaysThrowingFailTarget implements BiConsumer<String, String> {
int calls = 0;
@Override public void accept(String target, String reason) {
calls++;
throw new RuntimeException("boom");
}
}
} }