CB-580: fail a ticket when its member reaches a terminal health state
CI / build (pull_request) Successful in 54s
CI / contract (pull_request) Successful in 1m5s

FleetHealthMonitor now requires a failTarget BiConsumer<String,String>
collaborator (no defaulting overload) and calls it exactly once when a
member transitions into GONE or NEVER_READY, via CB-568's idempotent
target-wide abandon() operation. The reason string names the real
terminal state. failTarget invocation retries up to
MAX_FAIL_TARGET_ATTEMPTS (3) within the same transition if it throws,
and never refires on a later tick where the state is unchanged.

Bridged.java wires messages::abandon as the production failTarget.
This commit is contained in:
Dai Ha
2026-08-15 07:54:28 +02:00
parent 2f48e08f1f
commit 0af902ec43
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");
}
}
} }