From 3b2f395d3d8873e556ade0fae222d8a3eeef5725 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:31:51 +0200 Subject: [PATCH] M4: fail tickets on terminal health states --- .../src/main/java/dev/ltms/bridged/Bridged.java | 15 ++++++++++++++- .../ltms/bridged/health/FleetHealthMonitor.java | 14 +++++++++++++- 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index fbae9b7..2f50d0a 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -55,6 +55,7 @@ import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; @@ -254,6 +255,7 @@ public final class Bridged { // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied. Rendezvous rendezvous = new Rendezvous(); CompletionResolver completion = new CompletionResolver(agents, rendezvous); + AtomicReference messagesRef = new AtomicReference<>(); // CB-113: deliver only to an available worker (its MCP is connected), never its boot window. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY. MemberPresence presence = sessions.asPresence(); @@ -285,12 +287,14 @@ public final class Bridged { public void onTurnFailed(String target) { completion.onTurnFailed(target); sessions.onTurnFailed(target); + failTarget(messagesRef, target, "worker turn failed"); } @Override public void onTurnFailed(String target, String reason) { completion.onTurnFailed(target, reason); sessions.onTurnFailed(target); + failTarget(messagesRef, target, reason); } }; Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads), @@ -354,6 +358,7 @@ public final class Bridged { } MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, pushLoop, metrics); + messagesRef.set(messages); // Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller. final FleetHealthMonitor healthMonitor; @@ -361,7 +366,7 @@ public final class Bridged { Thread.ofVirtual().name("bridge-health-").unstarted(r)); if (cfg.health() != null && cfg.health().isEnabled()) { 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)) { @@ -485,6 +490,14 @@ public final class Bridged { return target -> presence.isPresent(target) || leads.get().containsKey(target); } + /** Fail outstanding tickets without changing the member lifecycle or worktree state. */ + private static void failTarget(AtomicReference messagesRef, String target, String reason) { + MessageService messages = messagesRef.get(); + if (messages != null) { + messages.abandon(target, reason); + } + } + /** * Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504). * 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..e0be638 100644 --- a/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java +++ b/bridged/src/main/java/dev/ltms/bridged/health/FleetHealthMonitor.java @@ -16,6 +16,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.function.LongSupplier; import java.util.function.Supplier; +import java.util.function.BiConsumer; /** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */ public final class FleetHealthMonitor { @@ -26,6 +27,7 @@ public final class FleetHealthMonitor { 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<>(); @@ -33,13 +35,20 @@ public final class FleetHealthMonitor { private static final boolean NOT_YET_OBSERVED = false; public FleetHealthMonitor(AgentControl agents, Supplier> roster, MessageService messages, - ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) { + ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) { + this(agents, roster, messages, scheduler, clock, intervalSeconds, (_, _) -> { }); + } + + public FleetHealthMonitor(AgentControl agents, Supplier> roster, MessageService messages, + 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 = failTarget; } /** Pure per-member decision seam. */ @@ -69,6 +78,9 @@ public final class FleetHealthMonitor { HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE), clock.getAsLong()); priors.put(session.terminalId(), decision.prior()); + if (decision.state() == HealthState.GONE || decision.state() == HealthState.NEVER_READY) { + failTarget.accept(session.terminalId(), "fleet health detected " + decision.state()); + } reportTransition(session.terminalId(), decision.state()); } priors.keySet().retainAll(current);