CB-580: fail a ticket when its member reaches a terminal health state #56

Merged
ltms merged 1 commits from worker/cb580-terminal-health-ed6058-21 into main 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 ->
Thread.ofVirtual().name("bridge-health-").unstarted(r));
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,
System::nanoTime, cfg.health().intervalOrDefault());
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
@@ -12,34 +12,50 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
/** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */
public final class FleetHealthMonitor {
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 Supplier<List<MemberSession>> roster;
private final MessageService messages;
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final long intervalSeconds;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = 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.
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,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) {
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
BiConsumer<String, String> failTarget) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.intervalSeconds = intervalSeconds;
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
}
/** Pure per-member decision seam. */
@@ -91,6 +107,34 @@ public final class FleetHealthMonitor {
} else if (previous != null && fault(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) {
@@ -20,8 +20,10 @@ import org.slf4j.LoggerFactory;
import java.util.Map;
import java.util.Set;
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.assertTrue;
class FleetHealthMonitorTest {
@Test void oneTickUsesOneFleetListForAnyRosterSize() {
@@ -37,7 +39,8 @@ class FleetHealthMonitorTest {
herdr.calls.clear();
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
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.stop();
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
@@ -49,7 +52,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60);
scheduler, () -> 1, 60, (_, _) -> { });
monitor.tick();
herdr.healthy(true);
monitor.tick();
@@ -68,7 +71,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
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.stop();
@@ -78,4 +81,90 @@ class FleetHealthMonitorTest {
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");
}
}
}