From 0af902ec438885a7e1d478d410672a01fc2791c2 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 07:54:28 +0200 Subject: [PATCH] CB-580: fail a ticket when its member reaches a terminal health state FleetHealthMonitor now requires a failTarget BiConsumer 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. --- .../main/java/dev/ltms/bridged/Bridged.java | 4 +- .../bridged/health/FleetHealthMonitor.java | 46 ++++++++- .../health/FleetHealthMonitorTest.java | 95 ++++++++++++++++++- 3 files changed, 140 insertions(+), 5 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 70a10d9..e38fa28 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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)) { diff --git a/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java index 79e4f84..f9ecd81 100644 --- a/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java +++ b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java @@ -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> roster; private final MessageService messages; private final ScheduledExecutorService scheduler; private final LongSupplier clock; private final long intervalSeconds; + private final BiConsumer failTarget; private final Map priors = new HashMap<>(); private final Map 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> roster, MessageService messages, - ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) { + ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds, + BiConsumer 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) { diff --git a/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java b/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java index 4d90b07..28c53e2 100644 --- a/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/health/FleetHealthMonitorTest.java @@ -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 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 { + final java.util.List 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 { + int calls = 0; + + @Override public void accept(String target, String reason) { + calls++; + throw new RuntimeException("boom"); + } + } } -- 2.52.0