M4: fail tickets on terminal health states
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 57s

This commit is contained in:
Dai Ha
2026-08-15 06:31:51 +02:00
parent 0edc6615fc
commit 3b2f395d3d
2 changed files with 27 additions and 2 deletions
@@ -55,6 +55,7 @@ import java.util.Map;
import java.util.Objects; import java.util.Objects;
import java.util.Set; import java.util.Set;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function; 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. // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
Rendezvous rendezvous = new Rendezvous(); Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = new CompletionResolver(agents, rendezvous); CompletionResolver completion = new CompletionResolver(agents, rendezvous);
AtomicReference<MessageService> messagesRef = new AtomicReference<>();
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window. // 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. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence(); MemberPresence presence = sessions.asPresence();
@@ -285,12 +287,14 @@ public final class Bridged {
public void onTurnFailed(String target) { public void onTurnFailed(String target) {
completion.onTurnFailed(target); completion.onTurnFailed(target);
sessions.onTurnFailed(target); sessions.onTurnFailed(target);
failTarget(messagesRef, target, "worker turn failed");
} }
@Override @Override
public void onTurnFailed(String target, String reason) { public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason); completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target); sessions.onTurnFailed(target);
failTarget(messagesRef, target, reason);
} }
}; };
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads), 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, MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics); pushLoop, metrics);
messagesRef.set(messages);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller. // Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
final FleetHealthMonitor healthMonitor; final FleetHealthMonitor healthMonitor;
@@ -361,7 +366,7 @@ public final class Bridged {
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()) {
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)) {
@@ -485,6 +490,14 @@ public final class Bridged {
return target -> presence.isPresent(target) || leads.get().containsKey(target); 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<MessageService> 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). * Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
* *
@@ -16,6 +16,7 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier; import java.util.function.LongSupplier;
import java.util.function.Supplier; import java.util.function.Supplier;
import java.util.function.BiConsumer;
/** 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 {
@@ -26,6 +27,7 @@ public final class FleetHealthMonitor {
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<>();
@@ -33,13 +35,20 @@ public final class FleetHealthMonitor {
private static final boolean NOT_YET_OBSERVED = false; private static final boolean NOT_YET_OBSERVED = false;
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) {
this(agents, roster, messages, scheduler, clock, intervalSeconds, (_, _) -> { });
}
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
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 = failTarget;
} }
/** Pure per-member decision seam. */ /** Pure per-member decision seam. */
@@ -69,6 +78,9 @@ public final class FleetHealthMonitor {
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE), HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
clock.getAsLong()); clock.getAsLong());
priors.put(session.terminalId(), decision.prior()); 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()); reportTransition(session.terminalId(), decision.state());
} }
priors.keySet().retainAll(current); priors.keySet().retainAll(current);