Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| bbf68f3e3c | |||
| fe2e5ede34 |
@@ -9,10 +9,8 @@ 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;
|
||||
@@ -70,12 +68,6 @@ 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;
|
||||
@@ -101,13 +93,6 @@ 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<>();
|
||||
|
||||
@@ -193,20 +178,7 @@ 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) {
|
||||
}
|
||||
|
||||
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) {
|
||||
String question, int nudgeCount) {
|
||||
}
|
||||
|
||||
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
|
||||
@@ -220,24 +192,6 @@ 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).
|
||||
@@ -282,15 +236,6 @@ 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.
|
||||
@@ -327,35 +272,17 @@ 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();
|
||||
boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty();
|
||||
boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) {
|
||||
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) {
|
||||
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;
|
||||
boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders;
|
||||
boolean unmappedTargetEligible = hasUnmappedTargetWork && unmappedTargetReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible && !unmappedTargetEligible) {
|
||||
if (!replyEligible && !ticketEligible && !questionEligible) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||
maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
@@ -467,48 +394,6 @@ 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. */
|
||||
@@ -540,25 +425,18 @@ 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);
|
||||
int incidentReminderCount = minIncidentNudgeCountFor(lead);
|
||||
int unmappedTargetReminderCount = minUnmappedTargetNudgeCountFor(lead);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
|
||||
incidentReminderCount, unmappedTargetReminderCount);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount,
|
||||
incidentReminderCount, unmappedTargetReminderCount);
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||
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, incidentsBefore,
|
||||
unmappedTargetsBefore);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -602,25 +480,11 @@ 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) {
|
||||
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) {
|
||||
Set<String> questionsBefore) {
|
||||
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))
|
||||
|| pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i))
|
||||
|| pendingUnmappedTargetKeysFor(lead).stream().anyMatch(i -> !unmappedTargetsBefore.contains(i));
|
||||
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t));
|
||||
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);
|
||||
@@ -631,22 +495,18 @@ 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 incidentReminderCount,
|
||||
int unmappedTargetReminderCount) {
|
||||
int questionReminderCount) {
|
||||
// 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);
|
||||
List<PendingIncident> incidents = pendingIncidentsFor(lead);
|
||||
List<PendingUnmappedTarget> unmappedTargets = pendingUnmappedTargetsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()
|
||||
&& unmappedTargets.isEmpty()) {
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) {
|
||||
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
||||
return;
|
||||
}
|
||||
String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets);
|
||||
String nudge = formatNudge(replyTargets, tickets, questions);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
||||
@@ -655,16 +515,6 @@ 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,
|
||||
@@ -676,13 +526,12 @@ 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, incidents, unmappedTargets);
|
||||
bumpNudgeCounts(replyTargets, tickets, questions);
|
||||
}
|
||||
|
||||
/** 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<PendingIncident> incidents,
|
||||
List<PendingUnmappedTarget> unmappedTargets) {
|
||||
List<PendingQuestion> questions) {
|
||||
for (String target : replyTargets) {
|
||||
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
||||
}
|
||||
@@ -695,14 +544,6 @@ 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. */
|
||||
@@ -715,8 +556,7 @@ 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<PendingIncident> incidents,
|
||||
List<PendingUnmappedTarget> unmappedTargets) {
|
||||
List<PendingQuestion> questions) {
|
||||
List<String> parts = new ArrayList<>();
|
||||
if (!replyTargets.isEmpty()) {
|
||||
parts.add(formatRepliesNudge(replyTargets));
|
||||
@@ -727,12 +567,6 @@ 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);
|
||||
}
|
||||
|
||||
@@ -772,16 +606,6 @@ 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 -----------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
@@ -803,10 +627,6 @@ public final class ReplyPushLoop {
|
||||
pendingReplies.clear();
|
||||
pendingTickets.clear();
|
||||
pendingQuestions.clear();
|
||||
pendingIncidents.clear();
|
||||
deliveredIncidents.clear();
|
||||
pendingUnmappedTargets.clear();
|
||||
deliveredUnmappedTargets.clear();
|
||||
}
|
||||
|
||||
/** @see #stop() */
|
||||
|
||||
@@ -46,7 +46,8 @@ public record MemberSession(
|
||||
String worktree,
|
||||
String branch,
|
||||
CharterReceipt charterReceipt,
|
||||
String agentSessionId) {
|
||||
String agentSessionId,
|
||||
String failureReason) {
|
||||
|
||||
/** One-shot worker lifecycle states. */
|
||||
public enum State {
|
||||
@@ -54,6 +55,7 @@ public record MemberSession(
|
||||
READY,
|
||||
BUSY,
|
||||
DONE,
|
||||
BACKEND_ERROR,
|
||||
FAILED,
|
||||
RELEASED
|
||||
}
|
||||
@@ -69,25 +71,35 @@ public record MemberSession(
|
||||
long lastActivityAtNanos, int turnCount, State state,
|
||||
String worktree, String branch) {
|
||||
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, null, null);
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, null, null, null);
|
||||
}
|
||||
|
||||
/** Backward-compatible shape without a backend failure reason. */
|
||||
public MemberSession(String paneId, String terminalId, String profile, MemberRole role,
|
||||
String cwd, String ownerTerminal, long spawnedAtNanos,
|
||||
long lastActivityAtNanos, int turnCount, State state,
|
||||
String worktree, String branch, CharterReceipt charterReceipt,
|
||||
String agentSessionId) {
|
||||
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, null);
|
||||
}
|
||||
|
||||
/** Return a copy of this session in {@code state}. */
|
||||
public MemberSession withState(State state) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
|
||||
}
|
||||
|
||||
/** Return a copy with {@code lastActivityAtNanos} updated to {@code nowNanos}. */
|
||||
public MemberSession withActivity(long nowNanos) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
nowNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
nowNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
|
||||
}
|
||||
|
||||
/** Return a copy with the turn count incremented and activity timestamped at {@code nowNanos}. */
|
||||
public MemberSession bumpTurn(long nowNanos) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -98,6 +110,12 @@ public record MemberSession(
|
||||
*/
|
||||
public MemberSession withAgentSessionId(String agentSessionId) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
|
||||
}
|
||||
|
||||
/** Return a copy with the durable backend failure detail. */
|
||||
public MemberSession withFailureReason(String failureReason) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -658,6 +658,9 @@ public final class SessionManager implements TurnListener {
|
||||
m.put("profile", session.profile());
|
||||
m.put("role", session.role() == null ? "dev" : session.role().wireName());
|
||||
m.put("state", session.state().name().toLowerCase());
|
||||
if (session.failureReason() != null) {
|
||||
m.put("failureReason", session.failureReason());
|
||||
}
|
||||
if (session.worktree() != null) {
|
||||
m.put("worktree", session.worktree());
|
||||
}
|
||||
@@ -718,6 +721,36 @@ public final class SessionManager implements TurnListener {
|
||||
completeTurn(target, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* Record that a member backend failed after a delegated ticket expired. This CAS loop accepts
|
||||
* either side of the completion race: {@code BUSY -> BACKEND_ERROR} or
|
||||
* {@code DONE -> BACKEND_ERROR}. A released or unknown session is never recreated.
|
||||
*
|
||||
* @return {@code true} when the error is recorded on a live session
|
||||
*/
|
||||
public boolean onBackendError(String target, String reason) {
|
||||
while (true) {
|
||||
MemberSession current = findByTerminal(target);
|
||||
if (current == null) {
|
||||
log.warn("backend error for unknown member terminal={}", target);
|
||||
return false;
|
||||
}
|
||||
if (current.state() == MemberSession.State.RELEASED) {
|
||||
return false;
|
||||
}
|
||||
if (current.state() == MemberSession.State.BACKEND_ERROR) {
|
||||
return true;
|
||||
}
|
||||
MemberSession updated = current.withState(MemberSession.State.BACKEND_ERROR)
|
||||
.withFailureReason(reason).withActivity(nowNanos.getAsLong());
|
||||
if (replace(current, updated)) {
|
||||
log.warn("member terminal={} pane={} transitioned {} -> BACKEND_ERROR: {}",
|
||||
target, current.paneId(), current.state(), reason);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasPostTurnAction(String target) {
|
||||
if (!clearAfterTurn) return false;
|
||||
@@ -736,10 +769,9 @@ public final class SessionManager implements TurnListener {
|
||||
if (current == null || current.state() != MemberSession.State.BUSY) return false;
|
||||
long now = nowNanos.getAsLong();
|
||||
MemberSession updated = current.withState(MemberSession.State.DONE).withActivity(now);
|
||||
if (replace(current, updated)) {
|
||||
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
|
||||
target, current.paneId(), updated.turnCount());
|
||||
}
|
||||
if (!replace(current, updated)) return false;
|
||||
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
|
||||
target, current.paneId(), updated.turnCount());
|
||||
if (contextCap > 0 && updated.turnCount() >= contextCap) {
|
||||
release(current.paneId());
|
||||
return false;
|
||||
|
||||
@@ -43,7 +43,6 @@ 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;
|
||||
@@ -530,106 +529,6 @@ 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
|
||||
@@ -1029,31 +928,6 @@ 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
|
||||
|
||||
@@ -300,6 +300,151 @@ class SessionManagerTest {
|
||||
assertTrue(sessions.roster().contains(updated), "FAILED is still in acquired-minus-released roster");
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendErrorWinsBeforeNormalCompletionFromBusy() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
|
||||
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
|
||||
sessions.onTurnComplete(session.terminalId());
|
||||
|
||||
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
|
||||
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state(),
|
||||
"normal completion must not restore DONE after a backend error from BUSY");
|
||||
assertEquals("backend exited", updated.failureReason());
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendErrorWinsAfterNormalCompletionFromDone() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
sessions.onTurnComplete(session.terminalId());
|
||||
|
||||
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
|
||||
|
||||
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
|
||||
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state(),
|
||||
"backend error must replace DONE when normal completion won first");
|
||||
assertEquals("backend exited", updated.failureReason());
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendErrorMemberCannotReceiveAnotherDelivery() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
|
||||
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
|
||||
assertEquals(MemberSession.State.BACKEND_ERROR, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"onDelivered must refuse a backend-error member");
|
||||
}
|
||||
|
||||
@Test
|
||||
void losingCompletionDoesNotReleaseOrClearABackendErrorMember() {
|
||||
assertLosingCompletionLeavesBackendErrorIntact(1, false,
|
||||
"a stale completion must not release a backend-error member at the context cap");
|
||||
assertLosingCompletionLeavesBackendErrorIntact(0, true,
|
||||
"a stale completion must not clear the context of a backend-error member");
|
||||
}
|
||||
|
||||
private void assertLosingCompletionLeavesBackendErrorIntact(int contextCap, boolean clearAfterTurn,
|
||||
String protectedAction) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
ClearContextSpyLauncher launcher = new ClearContextSpyLauncher(delegate);
|
||||
java.util.concurrent.atomic.AtomicReference<SessionManager> manager = new java.util.concurrent.atomic.AtomicReference<>();
|
||||
java.util.concurrent.atomic.AtomicReference<String> terminal = new java.util.concurrent.atomic.AtomicReference<>();
|
||||
java.util.concurrent.atomic.AtomicBoolean armed = new java.util.concurrent.atomic.AtomicBoolean();
|
||||
SessionManager sessions = new SessionManager(launcher, new RecordingWorktrees(), () -> {
|
||||
if (armed.compareAndSet(true, false)) {
|
||||
manager.get().onBackendError(terminal.get(), "backend exited during completion");
|
||||
}
|
||||
return 1;
|
||||
}, contextCap, clearAfterTurn);
|
||||
manager.set(sessions);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
terminal.set(session.terminalId());
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
armed.set(true);
|
||||
|
||||
assertFalse(sessions.onTurnCompleteWithPostAction(session.terminalId()),
|
||||
"the completion CAS loses after the injected backend error");
|
||||
|
||||
assertTrue(sessions.get(session.paneId()).isPresent(), protectedAction);
|
||||
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
|
||||
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state());
|
||||
assertEquals("backend exited during completion", updated.failureReason());
|
||||
assertFalse(herdr.called("pane.close"), protectedAction);
|
||||
assertEquals(0, launcher.clearContextCalls(), protectedAction);
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendErrorForUnknownTargetIsWarnedAndDoesNotCreateASession() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
try {
|
||||
SessionManager sessions = sessionManager(new FakeHerdr());
|
||||
|
||||
assertFalse(sessions.onBackendError("term_missing", "backend exited"));
|
||||
|
||||
assertTrue(sessions.roster().isEmpty(), "unknown target must not create a session");
|
||||
assertTrue(appender.list.stream().anyMatch(e -> e.getLevel().equals(Level.WARN)
|
||||
&& e.getFormattedMessage().contains("term_missing")),
|
||||
"unknown target is logged at WARN");
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendErrorForReleasedTargetDoesNotRecreateTheSession() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.release(session.paneId());
|
||||
|
||||
assertFalse(sessions.onBackendError(session.terminalId(), "backend exited"));
|
||||
assertTrue(sessions.get(session.paneId()).isEmpty(), "released session must stay absent");
|
||||
}
|
||||
|
||||
@Test
|
||||
void rosterViewIncludesFailureReasonOnlyForBackendError() {
|
||||
MemberSession failed = new MemberSession("pane-error", "term-error", "prof", MemberRole.DEV,
|
||||
"/cwd", null, 0, 0, 1, MemberSession.State.BACKEND_ERROR, null, null,
|
||||
null, null, "backend exited");
|
||||
MemberSession ordinary = new MemberSession("pane-ready", "term-ready", "prof", MemberRole.DEV,
|
||||
"/cwd", null, 0, 0, 0, MemberSession.State.READY, null, null);
|
||||
|
||||
Map<String, Object> failedView = SessionManager.rosterView(failed, null);
|
||||
Map<String, Object> ordinaryView = SessionManager.rosterView(ordinary, null);
|
||||
|
||||
assertEquals("backend_error", failedView.get("state"));
|
||||
assertEquals("backend exited", failedView.get("failureReason"));
|
||||
assertFalse(ordinaryView.containsKey("failureReason"),
|
||||
"ordinary rows must not gain a blank failureReason field");
|
||||
}
|
||||
|
||||
@Test
|
||||
void onTurnFailedIsLoggedAtWarnWithThePriorState() {
|
||||
// CB-564: this transition used to be a bare DEBUG "session marked failed" — a symptom with no
|
||||
@@ -1004,6 +1149,76 @@ class SessionManagerTest {
|
||||
}
|
||||
}
|
||||
|
||||
/** Delegates all launcher work while recording context clears. */
|
||||
private static final class ClearContextSpyLauncher implements PeerLauncher {
|
||||
private final PeerLauncher delegate;
|
||||
private int clearContextCalls;
|
||||
|
||||
ClearContextSpyLauncher(PeerLauncher delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return delegate.capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return delegate.capabilitiesFor(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return delegate.defaultProfile();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return delegate.effectiveCwd(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return delegate.parityOverlay(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return delegate.list();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return delegate.reapOrphanWorkers();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
delegate.stop(id);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
clearContextCalls++;
|
||||
return delegate.clearContext(id);
|
||||
}
|
||||
|
||||
int clearContextCalls() {
|
||||
return clearContextCalls;
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void acquireRecordsTheAgentSessionIdFromTheHandle() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
Reference in New Issue
Block a user