Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7840e9adf6 | |||
| ee932fd85b | |||
| 0f08b93659 | |||
| 432c1d92d1 | |||
| 60b7e67b42 | |||
| 3fbd43fe3f |
@@ -3,11 +3,13 @@ package dev.ltms.fleet.member;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.Tab;
|
||||
import dev.ltms.fleet.herdr.Workspace;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.StatusRefiner;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
@@ -151,6 +153,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
private final LongSupplier nowMillis; // monotonic clock (injectable for tests)
|
||||
private final Runnable sleeper; // sleep/wait hook (injectable for tests; never real-sleep in unit tests)
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 2: the same UNKNOWN-refinement {@code dev.ltms.fleet.inject.StatusPoller}
|
||||
* uses, reused here for the spawn-readiness gate. Constructed once from {@link #agents} — see
|
||||
* {@link #refinedInjectable(String, Agent)} for the corroboration that keeps it from firing on
|
||||
* a dead pane's bare shell prompt.
|
||||
*/
|
||||
private final StatusRefiner statusRefiner;
|
||||
|
||||
// Per-process token mixed into each peer name so a fresh process (nameSeq back at 0) cannot
|
||||
// collide with same-profile peers that outlived a restart. See startUniquelyNamed.
|
||||
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
|
||||
@@ -272,6 +282,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
this.memberCredentials = memberCredentials;
|
||||
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
|
||||
this.config = config;
|
||||
this.statusRefiner = new StatusRefiner(agents);
|
||||
}
|
||||
|
||||
// --- adapter seams -------------------------------------------------------------------------
|
||||
@@ -909,16 +920,35 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
// --- spawn-readiness gate (CB-306) ---------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Poll {@link AgentControl#status} until the pane reports an injectable state or the configured
|
||||
* Poll {@link AgentControl#get} until the pane reports an injectable state or the configured
|
||||
* timeout elapses. On timeout, close the pane (self-reap) and throw.
|
||||
*
|
||||
* <p>fleetd #176 fix 1: {@code agents.get} was previously unguarded here, so a herdr
|
||||
* {@code *_not_found} answer — which is what happens when the backend process EXITED rather
|
||||
* than being slow — propagated as a raw {@link HerdrException} instead of the
|
||||
* {@link PeerUnreachableException} every other failure path of this gate produces, and skipped
|
||||
* teardown ({@link #stop}) entirely, leaking the pane/tab. {@link #failFastOnGoneBackend} closes
|
||||
* that gap: it stops waiting immediately (never burns the rest of the timeout), runs the same
|
||||
* teardown the timeout path below runs, and throws with a message that says the backend exited
|
||||
* rather than that the pane was slow. Any other {@link HerdrException} still propagates
|
||||
* unchanged — this gate does not know how to recover from it.
|
||||
*/
|
||||
private void waitUntilInjectableOrThrow(String paneId) {
|
||||
long deadline = nowMillis.getAsLong() + spawnReadyTimeoutMs;
|
||||
Object lastStatus = null;
|
||||
long start = nowMillis.getAsLong();
|
||||
long deadline = start + spawnReadyTimeoutMs;
|
||||
AgentStatus lastStatus = null;
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
var status = agents.status(paneId);
|
||||
lastStatus = status;
|
||||
if (status.injectable()) {
|
||||
Agent sample;
|
||||
try {
|
||||
sample = agents.get(paneId);
|
||||
} catch (HerdrException e) {
|
||||
if (isAlreadyGone(e)) {
|
||||
failFastOnGoneBackend(paneId, e, nowMillis.getAsLong() - start);
|
||||
}
|
||||
throw e; // any other herdr failure is not ours to interpret — let it propagate
|
||||
}
|
||||
lastStatus = sample.status();
|
||||
if (lastStatus.injectable() || refinedInjectable(paneId, sample)) {
|
||||
log.debug("peer pane={} reached injectable state", paneId);
|
||||
return;
|
||||
}
|
||||
@@ -937,6 +967,66 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
+ spawnReadyTimeoutMs + "ms");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 1: the backend process exited while the gate was still waiting — herdr
|
||||
* answered {@code *_not_found} instead of ever reporting an injectable status. Fails
|
||||
* immediately (never burns the rest of {@link #spawnReadyTimeoutMs}), runs the exact same
|
||||
* teardown {@link #waitUntilInjectableOrThrow}'s timeout path runs, and throws a
|
||||
* {@link PeerUnreachableException} whose message says the process exited rather than that the
|
||||
* pane was slow.
|
||||
*
|
||||
* @throws PeerUnreachableException always — this method never returns normally
|
||||
*/
|
||||
private void failFastOnGoneBackend(String paneId, HerdrException cause, long elapsedMs) {
|
||||
String tail = readPaneQuietly(paneId); // read before stop() closes the pane
|
||||
log.warn("peer pane={} backend process exited after {}ms while waiting for injectable "
|
||||
+ "state (herdr: {}) — closing. Pane tail:\n{}",
|
||||
paneId, elapsedMs, cause.getMessage(), tail);
|
||||
stop(paneId);
|
||||
throw new PeerUnreachableException(
|
||||
"worker pane " + paneId + " backend process exited after " + elapsedMs
|
||||
+ "ms while waiting to become injectable (herdr reported: " + cause.getMessage()
|
||||
+ "). Pane tail:\n" + tail);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 2: resolve a raw {@link AgentStatus#UNKNOWN} sample into a trustworthy
|
||||
* injectable state via pane content, the same refinement {@code StatusPoller} applies during a
|
||||
* peer's working life — but corroborated, because the trap this gate is exposed to that the
|
||||
* poller is not: a pane whose backend has already exited settles at a plain shell prompt, and
|
||||
* that prompt commonly contains the same {@code ❯} glyph {@link StatusRefiner#classify} treats
|
||||
* as "idle at the Claude Code TUI prompt". Naively wiring the refiner in would turn "the backend
|
||||
* died" into "ready to inject" — strictly worse than today's timeout.
|
||||
*
|
||||
* <p>Two guards, both required:
|
||||
* <ul>
|
||||
* <li>{@link StatusRefiner#classify} is written for the Claude Code TUI only (its own javadoc
|
||||
* says so), so refinement only ever runs for the {@code claude} adapter — never for
|
||||
* {@code opencode} or any future non-Claude backend, whatever its pane looks like.</li>
|
||||
* <li>The refined status is accepted only when the <em>same</em> {@code agents.get} sample
|
||||
* ({@code sample}, taken once by the caller) still reports a non-null
|
||||
* {@link Agent#agentType()} — herdr's own "a supported backend is still detected here"
|
||||
* signal, read from the very sample the raw status came from so the two can never
|
||||
* disagree. A bare shell prompt left by an exited backend reports no {@code agentType}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Called only when {@code sample.status() == UNKNOWN}, so a healthy spawn — which never sees
|
||||
* {@code UNKNOWN} — triggers zero extra herdr calls; only a persistently-{@code unknown} pane
|
||||
* pays the one extra {@code agent.read} {@link StatusRefiner#refine} performs.
|
||||
*/
|
||||
private boolean refinedInjectable(String paneId, Agent sample) {
|
||||
if (sample.status() != AgentStatus.UNKNOWN) {
|
||||
return false;
|
||||
}
|
||||
if (!"claude".equals(namePrefix)) {
|
||||
return false; // StatusRefiner.classify reads a Claude Code TUI prompt specifically
|
||||
}
|
||||
if (sample.agentType() == null) {
|
||||
return false; // no corroborating liveness signal — could be a dead pane's bare shell
|
||||
}
|
||||
return statusRefiner.refine(paneId, AgentStatus.UNKNOWN, agents).injectable();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220: the pane's recent output, clipped, for the readiness-gate timeout log — or a
|
||||
* short note when it cannot be read. Best-effort by construction: this runs on a path that is
|
||||
|
||||
@@ -9,8 +9,10 @@ import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collection;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
@@ -68,6 +70,12 @@ public final class ReplyPushLoop {
|
||||
static final String QUESTIONS_NUDGE_FORMAT =
|
||||
"%d workers are paused on a question — run fleet_poll(ticket=...) for each, then answer "
|
||||
+ "with fleet_send(turnId=..., content=...): %s";
|
||||
static final String BACKEND_INCIDENT_NUDGE_FORMAT =
|
||||
"Backend credential %s is cooling for %d remaining seconds; affected profiles: %s; "
|
||||
+ "affected workers: %s. Run fleet_list to see more.";
|
||||
static final String BACKEND_TARGET_UNMAPPED_NUDGE_FORMAT =
|
||||
"A backend error on worker %s could not be mapped to a credential (%s) — no cool-off "
|
||||
+ "was applied. Run fleet_list to check that worker.";
|
||||
|
||||
private final PrimaryRegistry primaryRegistry;
|
||||
private final AgentControl agents;
|
||||
@@ -93,6 +101,13 @@ public final class ReplyPushLoop {
|
||||
* how depleted an older, still-open question's count is.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, PendingQuestion> pendingQuestions = new ConcurrentHashMap<>();
|
||||
/** Backend incidents awaiting one successful delivery, keyed by incident id and owning lead. */
|
||||
private final ConcurrentHashMap<IncidentLead, PendingIncident> pendingIncidents = new ConcurrentHashMap<>();
|
||||
/** Incident/lead pairs already delivered. They make repeated incident reports one-shot. */
|
||||
private final Set<IncidentLead> deliveredIncidents = ConcurrentHashMap.newKeySet();
|
||||
private final ConcurrentHashMap<UnmappedTargetLead, PendingUnmappedTarget> pendingUnmappedTargets =
|
||||
new ConcurrentHashMap<>();
|
||||
private final Set<UnmappedTargetLead> deliveredUnmappedTargets = ConcurrentHashMap.newKeySet();
|
||||
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */
|
||||
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
|
||||
|
||||
@@ -178,7 +193,20 @@ public final class ReplyPushLoop {
|
||||
* tracked per question, not per lead per source).
|
||||
*/
|
||||
private record PendingQuestion(String turnId, String ticket, String target, String lead,
|
||||
String question, int nudgeCount) {
|
||||
String question, int nudgeCount) {
|
||||
}
|
||||
|
||||
private record IncidentLead(String incidentId, String lead) {
|
||||
}
|
||||
|
||||
private record PendingIncident(IncidentLead key, String credential, List<String> profiles,
|
||||
List<String> targets, int remainingCoolOffSeconds, int nudgeCount) {
|
||||
}
|
||||
|
||||
private record UnmappedTargetLead(String target, String reason, String lead) {
|
||||
}
|
||||
|
||||
private record PendingUnmappedTarget(UnmappedTargetLead key, int nudgeCount) {
|
||||
}
|
||||
|
||||
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
|
||||
@@ -192,6 +220,24 @@ public final class ReplyPushLoop {
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
private List<PendingIncident> pendingIncidentsFor(String lead) {
|
||||
return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
|
||||
}
|
||||
|
||||
private Set<IncidentLead> pendingIncidentKeysFor(String lead) {
|
||||
return pendingIncidentsFor(lead).stream().map(PendingIncident::key)
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
private List<PendingUnmappedTarget> pendingUnmappedTargetsFor(String lead) {
|
||||
return pendingUnmappedTargets.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
|
||||
}
|
||||
|
||||
private Set<UnmappedTargetLead> pendingUnmappedTargetKeysFor(String lead) {
|
||||
return pendingUnmappedTargetsFor(lead).stream().map(PendingUnmappedTarget::key)
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
|
||||
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
|
||||
@@ -236,6 +282,15 @@ public final class ReplyPushLoop {
|
||||
return min == Integer.MAX_VALUE ? 0 : min;
|
||||
}
|
||||
|
||||
private int minIncidentNudgeCountFor(String lead) {
|
||||
return pendingIncidentsFor(lead).stream().mapToInt(PendingIncident::nudgeCount).min().orElse(0);
|
||||
}
|
||||
|
||||
private int minUnmappedTargetNudgeCountFor(String lead) {
|
||||
return pendingUnmappedTargetsFor(lead).stream().mapToInt(PendingUnmappedTarget::nudgeCount)
|
||||
.min().orElse(0);
|
||||
}
|
||||
|
||||
/**
|
||||
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
||||
* tickets alike — and return what the loop should do.
|
||||
@@ -272,17 +327,35 @@ public final class ReplyPushLoop {
|
||||
* @param questionReminderCount the lowest nudge count among questions open for this lead
|
||||
*/
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount) {
|
||||
return decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
|
||||
minIncidentNudgeCountFor(lead));
|
||||
}
|
||||
|
||||
/** As above, with backend incidents as a fourth, independently bounded source. */
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount,
|
||||
int incidentReminderCount) {
|
||||
return decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
|
||||
incidentReminderCount, minUnmappedTargetNudgeCountFor(lead));
|
||||
}
|
||||
|
||||
/** As above, with unmapped backend targets as a fifth, independently bounded source. */
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount,
|
||||
int incidentReminderCount, int unmappedTargetReminderCount) {
|
||||
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
||||
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
||||
boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) {
|
||||
boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty();
|
||||
boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) {
|
||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||
return Action.STOP;
|
||||
}
|
||||
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
||||
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
||||
boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible && !questionEligible) {
|
||||
boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders;
|
||||
boolean unmappedTargetEligible = hasUnmappedTargetWork && unmappedTargetReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible && !unmappedTargetEligible) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||
maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
@@ -394,6 +467,48 @@ public final class ReplyPushLoop {
|
||||
pendingQuestions.remove(turnId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Queue a one-shot backend credential outage notice for every distinct lead that owns an
|
||||
* affected worker. A successful injection records its {@code (incidentId, lead)} key, so a
|
||||
* repeat report never reminds that lead again.
|
||||
*/
|
||||
public void onBackendIncident(String incidentId, Collection<String> targets, String credential,
|
||||
Collection<String> profiles, int remainingCoolOffSeconds) {
|
||||
Map<String, List<String>> targetsByLead = new ConcurrentHashMap<>();
|
||||
for (String target : targets) {
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.warn("push: backend incident {} has no known lead for target {}", incidentId, target);
|
||||
continue;
|
||||
}
|
||||
targetsByLead.computeIfAbsent(lead.get(), _ -> new ArrayList<>()).add(target);
|
||||
}
|
||||
List<String> profileNames = profiles.stream().sorted().toList();
|
||||
for (var entry : targetsByLead.entrySet()) {
|
||||
IncidentLead key = new IncidentLead(incidentId, entry.getKey());
|
||||
if (deliveredIncidents.contains(key)) continue;
|
||||
pendingIncidents.putIfAbsent(key, new PendingIncident(key, credential, profileNames,
|
||||
entry.getValue().stream().sorted().toList(), remainingCoolOffSeconds, 0));
|
||||
startOrCoalesce(entry.getKey());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Tell the owning lead that a classified backend error could not be tied to a credential.
|
||||
* Without an owning lead, emit a warning because no control can act on the target.
|
||||
*/
|
||||
public void onBackendTargetUnmapped(String target, String reason) {
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.warn("push: backend target {} could not map to a credential: {}", target, reason);
|
||||
return;
|
||||
}
|
||||
UnmappedTargetLead key = new UnmappedTargetLead(target, reason, lead.get());
|
||||
if (deliveredUnmappedTargets.contains(key)) return;
|
||||
pendingUnmappedTargets.putIfAbsent(key, new PendingUnmappedTarget(key, 0));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
// --- the schedule ----------------------------------------------------------------------------
|
||||
|
||||
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
||||
@@ -425,18 +540,25 @@ public final class ReplyPushLoop {
|
||||
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
||||
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
||||
Set<String> questionsBefore = pendingQuestionTurnIdsFor(lead);
|
||||
Set<IncidentLead> incidentsBefore = pendingIncidentKeysFor(lead);
|
||||
Set<UnmappedTargetLead> unmappedTargetsBefore = pendingUnmappedTargetKeysFor(lead);
|
||||
int replyReminderCount = minReplyNudgeCountFor(lead);
|
||||
int ticketReminderCount = minTicketNudgeCountFor(lead);
|
||||
int questionReminderCount = minQuestionNudgeCountFor(lead);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||
int incidentReminderCount = minIncidentNudgeCountFor(lead);
|
||||
int unmappedTargetReminderCount = minUnmappedTargetNudgeCountFor(lead);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
|
||||
incidentReminderCount, unmappedTargetReminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
|
||||
incidentReminderCount, unmappedTargetReminderCount);
|
||||
scheduleNext(lead);
|
||||
}
|
||||
// Re-check after the configured backoff; the lead may become injectable soon.
|
||||
case WAIT_BUSY -> scheduleNext(lead);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, incidentsBefore,
|
||||
unmappedTargetsBefore);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -480,11 +602,25 @@ public final class ReplyPushLoop {
|
||||
* decision-to-release window reclaims the schedule slot exactly like a raced-in reply or ticket.
|
||||
*/
|
||||
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
|
||||
Set<String> questionsBefore) {
|
||||
Set<String> questionsBefore) {
|
||||
stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, pendingIncidentKeysFor(lead));
|
||||
}
|
||||
|
||||
private void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
|
||||
Set<String> questionsBefore, Set<IncidentLead> incidentsBefore) {
|
||||
stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, incidentsBefore,
|
||||
pendingUnmappedTargetKeysFor(lead));
|
||||
}
|
||||
|
||||
private void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
|
||||
Set<String> questionsBefore, Set<IncidentLead> incidentsBefore,
|
||||
Set<UnmappedTargetLead> unmappedTargetsBefore) {
|
||||
activeLeads.remove(lead);
|
||||
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|
||||
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t))
|
||||
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t));
|
||||
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t))
|
||||
|| pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i))
|
||||
|| pendingUnmappedTargetKeysFor(lead).stream().anyMatch(i -> !unmappedTargetsBefore.contains(i));
|
||||
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
||||
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
||||
scheduleNext(lead);
|
||||
@@ -495,18 +631,22 @@ public final class ReplyPushLoop {
|
||||
|
||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount,
|
||||
int questionReminderCount) {
|
||||
int questionReminderCount, int incidentReminderCount,
|
||||
int unmappedTargetReminderCount) {
|
||||
// Re-read rather than threading it down from decide(): a reply can drain, a ticket be
|
||||
// collected, or a question be answered (or another arrive), between the decision and the
|
||||
// injection.
|
||||
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
||||
List<PendingTicket> tickets = pendingTicketsFor(lead);
|
||||
List<PendingQuestion> questions = pendingQuestionsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) {
|
||||
List<PendingIncident> incidents = pendingIncidentsFor(lead);
|
||||
List<PendingUnmappedTarget> unmappedTargets = pendingUnmappedTargetsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()
|
||||
&& unmappedTargets.isEmpty()) {
|
||||
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
||||
return;
|
||||
}
|
||||
String nudge = formatNudge(replyTargets, tickets, questions);
|
||||
String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
||||
@@ -515,6 +655,16 @@ public final class ReplyPushLoop {
|
||||
questionReminderCount + 1, maxReminders,
|
||||
replyTargets.size(), tickets.size(), questions.size());
|
||||
countNudge("delivered");
|
||||
for (PendingIncident incident : incidents) {
|
||||
if (pendingIncidents.remove(incident.key(), incident)) {
|
||||
deliveredIncidents.add(incident.key());
|
||||
}
|
||||
}
|
||||
for (PendingUnmappedTarget unmappedTarget : unmappedTargets) {
|
||||
if (pendingUnmappedTargets.remove(unmappedTarget.key(), unmappedTarget)) {
|
||||
deliveredUnmappedTargets.add(unmappedTarget.key());
|
||||
}
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
@@ -526,12 +676,13 @@ public final class ReplyPushLoop {
|
||||
// item already at or over the cap keeps riding along in the text (still pending, still
|
||||
// named) but its extra bumps here are inert: decide() already treats it as ineligible once
|
||||
// its count reaches maxReminders.
|
||||
bumpNudgeCounts(replyTargets, tickets, questions);
|
||||
bumpNudgeCounts(replyTargets, tickets, questions, incidents, unmappedTargets);
|
||||
}
|
||||
|
||||
/** Record that every one of these items was just named in a sent (or attempted) nudge. */
|
||||
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets,
|
||||
List<PendingQuestion> questions) {
|
||||
List<PendingQuestion> questions, List<PendingIncident> incidents,
|
||||
List<PendingUnmappedTarget> unmappedTargets) {
|
||||
for (String target : replyTargets) {
|
||||
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
||||
}
|
||||
@@ -544,6 +695,14 @@ public final class ReplyPushLoop {
|
||||
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
|
||||
e.nudgeCount() + 1));
|
||||
}
|
||||
for (PendingIncident incident : incidents) {
|
||||
pendingIncidents.computeIfPresent(incident.key(), (id, e) -> new PendingIncident(e.key(),
|
||||
e.credential(), e.profiles(), e.targets(), e.remainingCoolOffSeconds(), e.nudgeCount() + 1));
|
||||
}
|
||||
for (PendingUnmappedTarget unmappedTarget : unmappedTargets) {
|
||||
pendingUnmappedTargets.computeIfPresent(unmappedTarget.key(), (id, e) ->
|
||||
new PendingUnmappedTarget(e.key(), e.nudgeCount() + 1));
|
||||
}
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
@@ -556,7 +715,8 @@ public final class ReplyPushLoop {
|
||||
|
||||
/** Render everything pending for one lead as a single nudge line. */
|
||||
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets,
|
||||
List<PendingQuestion> questions) {
|
||||
List<PendingQuestion> questions, List<PendingIncident> incidents,
|
||||
List<PendingUnmappedTarget> unmappedTargets) {
|
||||
List<String> parts = new ArrayList<>();
|
||||
if (!replyTargets.isEmpty()) {
|
||||
parts.add(formatRepliesNudge(replyTargets));
|
||||
@@ -567,6 +727,12 @@ public final class ReplyPushLoop {
|
||||
if (!questions.isEmpty()) {
|
||||
parts.add(formatQuestionsNudge(questions));
|
||||
}
|
||||
if (!incidents.isEmpty()) {
|
||||
parts.addAll(incidents.stream().map(ReplyPushLoop::formatBackendIncidentNudge).toList());
|
||||
}
|
||||
if (!unmappedTargets.isEmpty()) {
|
||||
parts.addAll(unmappedTargets.stream().map(ReplyPushLoop::formatUnmappedTargetNudge).toList());
|
||||
}
|
||||
return String.join(" | ", parts);
|
||||
}
|
||||
|
||||
@@ -606,6 +772,16 @@ public final class ReplyPushLoop {
|
||||
return QUESTIONS_NUDGE_FORMAT.formatted(pending.size(), ids);
|
||||
}
|
||||
|
||||
private static String formatBackendIncidentNudge(PendingIncident incident) {
|
||||
return BACKEND_INCIDENT_NUDGE_FORMAT.formatted(incident.credential(), incident.remainingCoolOffSeconds(),
|
||||
String.join(", ", incident.profiles()), String.join(", ", incident.targets()));
|
||||
}
|
||||
|
||||
private static String formatUnmappedTargetNudge(PendingUnmappedTarget unmappedTarget) {
|
||||
return BACKEND_TARGET_UNMAPPED_NUDGE_FORMAT.formatted(unmappedTarget.key().target(),
|
||||
unmappedTarget.key().reason());
|
||||
}
|
||||
|
||||
// --- lifecycle -----------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
@@ -627,6 +803,10 @@ public final class ReplyPushLoop {
|
||||
pendingReplies.clear();
|
||||
pendingTickets.clear();
|
||||
pendingQuestions.clear();
|
||||
pendingIncidents.clear();
|
||||
deliveredIncidents.clear();
|
||||
pendingUnmappedTargets.clear();
|
||||
deliveredUnmappedTargets.clear();
|
||||
}
|
||||
|
||||
/** @see #stop() */
|
||||
|
||||
@@ -44,10 +44,13 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private String agentSendErrorCode = null;
|
||||
private boolean noPanes = false;
|
||||
private volatile String agentStatus = "idle"; // steady-state agent.get status
|
||||
private volatile String agentType = "claude"; // detected agent kind on agent.get; null = undetected
|
||||
private volatile String readText = "worker transcript tail"; // canned agent.read output
|
||||
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
|
||||
private String pinnedStartTerminal;
|
||||
private String pinnedStartPane;
|
||||
private volatile int agentGetOkCalls = Integer.MAX_VALUE; // how many agent.get calls succeed first
|
||||
private volatile String agentGetFailCode = null; // error code every agent.get call after that reports
|
||||
|
||||
public FakeHerdr healthy(boolean h) {
|
||||
this.healthy = h;
|
||||
@@ -103,6 +106,27 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the detected agent kind ({@code "agent"} field) that {@code agent.get} reports —
|
||||
* {@code null} models herdr not (or no longer) detecting a supported backend in the pane, e.g.
|
||||
* a bare shell prompt (fleetd #176 fix 2 corroboration test).
|
||||
*/
|
||||
public FakeHerdr agentType(String type) {
|
||||
this.agentType = type;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code agent.get} succeed normally for its first {@code okCalls} invocations, then fail
|
||||
* every call after that with {@code code} — fleetd #176 fix 1's "backend exited mid-wait"
|
||||
* fixture. {@code okCalls == 0} fails from the very first call.
|
||||
*/
|
||||
public FakeHerdr agentGetFailsWithAfter(int okCalls, String code) {
|
||||
this.agentGetOkCalls = okCalls;
|
||||
this.agentGetFailCode = code;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** The text {@code agent.read} returns (the CB-106 completion scrape). */
|
||||
public FakeHerdr readText(String text) {
|
||||
this.readText = text;
|
||||
@@ -208,10 +232,21 @@ public final class FakeHerdr implements HerdrClient {
|
||||
}
|
||||
yield mapper.readTree("{\"type\":\"ok\"}");
|
||||
}
|
||||
case "agent.get" -> mapper.readTree(("""
|
||||
{"type":"agent_info","agent":{"terminal_id":"term_a","agent":"claude",
|
||||
case "agent.get" -> {
|
||||
if (agentGetFailCode != null) {
|
||||
long getCalls = calls.stream().filter(c -> c.method().equals("agent.get")).count();
|
||||
if (getCalls > agentGetOkCalls) {
|
||||
throw new HerdrException(
|
||||
"herdr error [" + agentGetFailCode + "]: agent target not found",
|
||||
agentGetFailCode, null);
|
||||
}
|
||||
}
|
||||
String agentField = agentType == null ? "null" : "\"" + agentType + "\"";
|
||||
yield mapper.readTree(("""
|
||||
{"type":"agent_info","agent":{"terminal_id":"term_a","agent":%s,
|
||||
"agent_status":"%s","workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"}}""")
|
||||
.formatted(agentStatus));
|
||||
.formatted(agentField, agentStatus));
|
||||
}
|
||||
case "agent.read" -> mapper.readTree(mapper.writeValueAsString(
|
||||
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText))));
|
||||
case "agent.start" -> {
|
||||
|
||||
@@ -963,6 +963,152 @@ class ClaudeCodeLauncherTest {
|
||||
assertDoesNotThrow(() -> UUID.fromString(handle.id()));
|
||||
}
|
||||
|
||||
// --- fleetd #176 fix 1: fail fast when the backend process exits mid-wait -------------------
|
||||
|
||||
@Test
|
||||
void spawnFailsFastWhenBackendProcessExitsMidWaitInsteadOfBurningTheTimeout() {
|
||||
// First agent.get sees UNKNOWN (one normal tick); the second reports the pane gone, exactly
|
||||
// what herdr answers when the backend process has already exited. The timeout is generous
|
||||
// (60s) so a test that wrongly falls through to the old unguarded call — and therefore
|
||||
// waits out the whole window — is unambiguously distinguishable from one that fails fast.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown");
|
||||
herdr.agentGetFailsWithAfter(1, "pane_not_found");
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
60_000, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
PeerUnreachableException ex = assertThrows(
|
||||
PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertTrue(ex.getMessage().contains("w9:pRoot_1"),
|
||||
"exception message references the paneId: " + ex.getMessage());
|
||||
assertTrue(ex.getMessage().toLowerCase().contains("exited"),
|
||||
"exception message says the backend exited, not that the pane was slow: "
|
||||
+ ex.getMessage());
|
||||
assertTrue(clock[0] < 60_000,
|
||||
"the gate must not burn the rest of the 60s timeout: clock only reached " + clock[0]);
|
||||
long getCalls = herdr.calls.stream().filter(c -> c.method().equals("agent.get")).count();
|
||||
assertEquals(2, getCalls,
|
||||
"exactly one normal poll then the not_found answer — no further polling after that: "
|
||||
+ getCalls);
|
||||
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"the pane is torn down (no orphan) even on the fail-fast path");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnLetsAnUnrelatedHerdrErrorPropagateUnchanged() {
|
||||
// Fix 1 must only special-case a "*_not_found" answer. Any other herdr failure keeps
|
||||
// propagating as-is — this gate does not know how to recover from it.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown");
|
||||
herdr.agentGetFailsWithAfter(0, "internal_error");
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
60_000, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
dev.ltms.fleet.herdr.HerdrException ex = assertThrows(
|
||||
dev.ltms.fleet.herdr.HerdrException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertEquals("internal_error", ex.code());
|
||||
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"an error this gate does not recognize is not this gate's teardown to run");
|
||||
}
|
||||
|
||||
// --- fleetd #176 fix 2: corroborated UNKNOWN refinement --------------------------------------
|
||||
|
||||
@Test
|
||||
void refinedIdleIsNotAcceptedWhenAgentTypeIsNull() {
|
||||
// The trap fix 2 must close: a pane sitting at a bare shell prompt after its backend exited
|
||||
// still contains the same "❯" glyph StatusRefiner.classify treats as "idle at the Claude
|
||||
// Code TUI prompt". Without the agentType corroboration this would be misread as injectable
|
||||
// and the gate would hand back a peer that never started. herdr reports no agentType for
|
||||
// that bare shell.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown"); // always UNKNOWN
|
||||
herdr.agentType(null); // no supported backend detected — could be a bare shell
|
||||
herdr.readText("some-host:~ user$ ❯ "); // looks exactly like an idle Claude Code prompt
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
500, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
PeerUnreachableException ex = assertThrows(
|
||||
PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertTrue(clock[0] >= 500,
|
||||
"a null agentType must not let the '❯' prompt refine to injectable — the gate has "
|
||||
+ "to wait out the full timeout: clock only reached " + clock[0]);
|
||||
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"the pane is torn down on timeout, same as any other never-injectable spawn");
|
||||
}
|
||||
|
||||
@Test
|
||||
void refinedIdleIsAcceptedWhenAgentTypeCorroboratesLiveness() {
|
||||
// The positive case: a genuinely live claude pane that herdr misreports as UNKNOWN (CB-115)
|
||||
// still resolves to injectable once agentType corroborates that a supported backend is
|
||||
// detected in the same sample.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown"); // always UNKNOWN at the raw level
|
||||
herdr.agentType("claude"); // herdr still detects a live claude backend
|
||||
herdr.readText("? for shortcuts"); // StatusRefiner.classify's idle footer marker
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
60_000, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
||||
|
||||
assertNotNull(handle, "a corroborated refined-IDLE sample lets the spawn succeed");
|
||||
assertTrue(clock[0] < 60_000,
|
||||
"refinement must resolve well before the timeout: clock reached " + clock[0]);
|
||||
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"no pane close — the peer is genuinely injectable");
|
||||
}
|
||||
|
||||
@Test
|
||||
void refinementNeverFiresWhenRawStatusIsAlreadyInjectable() {
|
||||
// Acceptance criterion 4: refinement must trigger only on a raw UNKNOWN sample. Seed the
|
||||
// pane content with an active-generation marker that StatusRefiner.classify would read as
|
||||
// WORKING (never injectable) if refine() were wrongly invoked here — proving that a raw
|
||||
// IDLE status short-circuits before refine() (and its extra agent.read call) ever runs.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("idle"); // already injectable at the raw level
|
||||
herdr.agentType("claude");
|
||||
herdr.readText("esc to interrupt"); // would classify as WORKING if refine() ran anyway
|
||||
|
||||
long[] clock = {0};
|
||||
ClaudeCodeLauncher gated = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
5000, () -> clock[0], () -> clock[0] += 300);
|
||||
|
||||
PeerHandle handle = gated.spawn(new SpawnRequest(null, null, null));
|
||||
|
||||
assertNotNull(handle, "an already-injectable raw status succeeds without ever refining");
|
||||
assertFalse(herdr.called("agent.read"),
|
||||
"refine() must never run (and so never call agent.read) when the raw status is "
|
||||
+ "already injectable");
|
||||
}
|
||||
|
||||
// --- CB-511: worker environment seeding -----------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -396,6 +396,46 @@ class OpenCodeLauncherTest {
|
||||
assertNotNull(ex.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 2's per-adapter guard: {@code StatusRefiner.classify} is written for the
|
||||
* Claude Code TUI only, and the spawn-readiness gate must never run it against an opencode
|
||||
* pane. {@code agentType("opencode")} deliberately satisfies the OTHER guard (the corroborating
|
||||
* liveness check) so it cannot be what makes this test pass — only the {@code namePrefix}
|
||||
* check can be. If that check were ever removed, this pane's {@code ❯} content would refine
|
||||
* straight to IDLE and the gate would report ready before the backend actually was.
|
||||
*/
|
||||
@Test
|
||||
void opencodePaneIsNeverRefinedEvenWhenItsContentLooksLikeAnIdleClaudePrompt(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown"); // always UNKNOWN
|
||||
herdr.agentType("opencode"); // non-null — satisfies the liveness guard on its own
|
||||
herdr.readText("some-host:~ user$ ❯ "); // content StatusRefiner.classify reads as IDLE
|
||||
long[] clock = {0};
|
||||
OpenCodeLauncher svc = new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("gemini", opencodeCfg(null, null, null)), "gemini", _ -> null,
|
||||
1000, () -> clock[0], () -> clock[0] += 50, root, root);
|
||||
|
||||
assertThrows(PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertTrue(clock[0] >= 1000,
|
||||
"an opencode pane must never refine to injectable, however its content reads — "
|
||||
+ "the gate has to wait out the full timeout: clock only reached " + clock[0]);
|
||||
// A blanket "agent.read is never called" does not hold here: the timeout path itself reads
|
||||
// the pane tail (source=recent) for its own log message, on every timeout, regardless of
|
||||
// adapter — see HerdrPeerLauncher.readPaneQuietly. So assert on the refiner's OWN probe
|
||||
// source (StatusRefiner.PROBE_SOURCE = "detection") instead — that call happens only inside
|
||||
// StatusRefiner.refine, so its absence proves the refiner itself was never reached for this
|
||||
// opencode pane, not merely that its answer was discarded.
|
||||
long detectionReads = herdr.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.read"))
|
||||
.filter(c -> "detection".equals(((Map<?, ?>) c.params()).get("source")))
|
||||
.count();
|
||||
assertEquals(0, detectionReads,
|
||||
"the refiner's own pane probe (source=detection) must never run against an "
|
||||
+ "opencode pane — the namePrefix guard has to stop it before that call");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnReturnsHandleWhenGateDisabled(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -43,6 +43,7 @@ class ReplyPushLoopTest {
|
||||
private static final String PRIMARY = "term_primary";
|
||||
private static final String WORKER = "term_worker";
|
||||
private static final String WORKER2 = "term_worker2";
|
||||
private static final String OTHER_PRIMARY = "term_other_primary";
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
|
||||
private PrimaryRegistry registry;
|
||||
@@ -529,6 +530,106 @@ class ReplyPushLoopTest {
|
||||
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
|
||||
}
|
||||
|
||||
// --- CB-201: backend credential outage nudges -----------------------------------------------
|
||||
|
||||
@Test
|
||||
void oneBackendIncidentForTwoWorkersOfOneLeadSendsOnceAndIsOneShot() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
registry.recordDelegation(WORKER, PRIMARY);
|
||||
registry.recordDelegation(WORKER2, PRIMARY);
|
||||
var loop = loop(2, 100);
|
||||
|
||||
loop.onBackendIncident("incident-1", List.of(WORKER, WORKER2), "cred-a",
|
||||
List.of("terra", "codex"), 47);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one lead gets one outage injection");
|
||||
Thread.sleep(200);
|
||||
assertEquals(1, rec.sendCount(), "two workers for one lead must send only one outage notice");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("cred-a"));
|
||||
assertTrue(nudge.contains("codex, terra"));
|
||||
assertTrue(nudge.contains(WORKER) && nudge.contains(WORKER2));
|
||||
assertTrue(nudge.contains("47 remaining seconds"));
|
||||
assertTrue(nudge.contains("fleet_list"));
|
||||
|
||||
loop.onBackendIncident("incident-1", List.of(WORKER, WORKER2), "cred-a",
|
||||
List.of("terra", "codex"), 47);
|
||||
Thread.sleep(250);
|
||||
assertEquals(1, rec.sendCount(), "a delivered incident id must never be sent again");
|
||||
}
|
||||
|
||||
@Test
|
||||
void oneBackendIncidentSendsOnceToEachAffectedLead() throws Exception {
|
||||
var rec = recordingClient();
|
||||
rec.sendLatch = new CountDownLatch(2);
|
||||
agents = new AgentControl(rec);
|
||||
registry.recordDelegation(WORKER, PRIMARY);
|
||||
registry.recordDelegation(WORKER2, OTHER_PRIMARY);
|
||||
var loop = loop(1, 50);
|
||||
|
||||
loop.onBackendIncident("incident-2", List.of(WORKER, WORKER2), "cred-a", List.of("terra"), 60);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "each affected lead gets its one notice");
|
||||
assertEquals(2, rec.sendCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendIncidentWaitsForBusyLeadThenCombinesWithFailedTicket() throws Exception {
|
||||
var rec = new BusyThenIdleHerdrClient(2);
|
||||
agents = new AgentControl(rec);
|
||||
var loop = loop(1, 50);
|
||||
|
||||
loop.onTicketTerminal("task-failed", WORKER, true);
|
||||
loop.onBackendIncident("incident-3", List.of(WORKER), "cred-b", List.of("terra"), 31);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "busy lead should receive deferred combined nudge");
|
||||
assertEquals(1, rec.sendCount(), "ticket and outage must use the same injection");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("task-failed") && nudge.contains("cred-b"),
|
||||
"the one nudge must name both the failed ticket and outage: " + nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void failedBackendIncidentSendRetriesWithinTheExistingReminderCap() throws Exception {
|
||||
var rec = new ThrowOnceHerdrClient();
|
||||
agents = new AgentControl(rec);
|
||||
var loop = loop(2, 50);
|
||||
|
||||
loop.onBackendIncident("incident-4", List.of(WORKER), "cred-c", List.of("terra"), 20);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "failed first send must retry under maxReminders");
|
||||
assertEquals(2, rec.promptAttempts.get(), "one failed send and one retry use the existing cap");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unmappedBackendTargetUsesTheKnownLeadSchedule() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
var loop = loop(1, 50);
|
||||
|
||||
loop.onBackendTargetUnmapped(WORKER, "no credential matched");
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
|
||||
String nudge = ((Map<?, ?>) rec.sentParams().getFirst().getValue()).get("text").toString();
|
||||
assertEquals("A backend error on worker term_worker could not be mapped to a credential "
|
||||
+ "(no credential matched) — no cool-off was applied. "
|
||||
+ "Run fleet_list to check that worker.",
|
||||
nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopClearsPendingBackendIncidents() {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(1, 100_000);
|
||||
loop.onBackendIncident("incident-stop", List.of(WORKER), "cred-stop", List.of("terra"), 60);
|
||||
|
||||
loop.stop();
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0, 0),
|
||||
"stop must clear a pending backend incident as well as older sources");
|
||||
}
|
||||
|
||||
// --- CB-590 follow-up: per-source reminder budgets — the regression this round exists for ---
|
||||
|
||||
@Test
|
||||
@@ -928,6 +1029,31 @@ class ReplyPushLoopTest {
|
||||
return new RecordingHerdrClient();
|
||||
}
|
||||
|
||||
/** Fails its first prompt call, then records the retry. */
|
||||
private static final class ThrowOnceHerdrClient implements HerdrClient {
|
||||
private final AtomicInteger promptAttempts = new AtomicInteger();
|
||||
private final CountDownLatch sendLatch = new CountDownLatch(1);
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", PRIMARY).put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.prompt".equals(method) && promptAttempts.getAndIncrement() == 0) {
|
||||
throw new IllegalStateException("first prompt fails");
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
sendLatch.countDown();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Thread-safe fake that reports {@code working} (not injectable) for its first
|
||||
* {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
|
||||
|
||||
Reference in New Issue
Block a user