M4: fail tickets on terminal health states
This commit is contained in:
@@ -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);
|
||||||
|
|||||||
Reference in New Issue
Block a user