From ee932fd85bd8fd9f19eaa76a4915a95b952d3745 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:28:51 +0700 Subject: [PATCH 1/2] CB-201: nudge leads about backend outages --- .../dev/ltms/fleet/msg/ReplyPushLoop.java | 136 ++++++++++++++++-- .../dev/ltms/fleet/msg/ReplyPushLoopTest.java | 123 ++++++++++++++++ 2 files changed, 245 insertions(+), 14 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java index 480751f..c85228f 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -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,9 @@ 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."; private final PrimaryRegistry primaryRegistry; private final AgentControl agents; @@ -93,6 +98,10 @@ public final class ReplyPushLoop { * how depleted an older, still-open question's count is. */ private final ConcurrentHashMap pendingQuestions = new ConcurrentHashMap<>(); + /** Backend incidents awaiting one successful delivery, keyed by incident id and owning lead. */ + private final ConcurrentHashMap pendingIncidents = new ConcurrentHashMap<>(); + /** Incident/lead pairs already delivered. They make repeated incident reports one-shot. */ + private final Set deliveredIncidents = ConcurrentHashMap.newKeySet(); /** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */ private final ConcurrentHashMap activeLeads = new ConcurrentHashMap<>(); @@ -178,7 +187,14 @@ 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 profiles, + List targets, int remainingCoolOffSeconds, int nudgeCount) { } /** Questions still open for {@code lead}, snapshotted fresh for one tick. */ @@ -192,6 +208,15 @@ public final class ReplyPushLoop { .collect(Collectors.toUnmodifiableSet()); } + private List pendingIncidentsFor(String lead) { + return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList(); + } + + private Set pendingIncidentKeysFor(String lead) { + return pendingIncidentsFor(lead).stream().map(PendingIncident::key) + .collect(Collectors.toUnmodifiableSet()); + } + /** * The reply-source reminder count {@link #decide} should see for {@code lead} on this tick: * the minimum nudge count among the reply targets currently pending for it (CB-598). @@ -236,6 +261,10 @@ 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); + } + /** * Pure decision function: examine everything pending for {@code lead} — reply targets and * tickets alike — and return what the loop should do. @@ -272,17 +301,26 @@ 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) { boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty(); boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty(); boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty(); - if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) { + boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty(); + if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork) { 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; + if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible) { log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping", maxReminders, lead); countNudge("exhausted"); @@ -394,6 +432,46 @@ 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 targets, String credential, + Collection profiles, int remainingCoolOffSeconds) { + Map> 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 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; + } + onBackendIncident("unmapped:" + target + ":" + reason, List.of(target), "unknown", + List.of("unmapped backend error: " + reason), 0); + } + // --- the schedule ---------------------------------------------------------------------------- /** Start a reminder schedule for {@code lead}, or join the one already running. */ @@ -425,18 +503,22 @@ public final class ReplyPushLoop { Set repliesBefore = pendingReplyTargetsFor(lead); Set ticketsBefore = pendingTicketIdsFor(lead); Set questionsBefore = pendingQuestionTurnIdsFor(lead); + Set incidentsBefore = pendingIncidentKeysFor(lead); int replyReminderCount = minReplyNudgeCountFor(lead); int ticketReminderCount = minTicketNudgeCountFor(lead); int questionReminderCount = minQuestionNudgeCountFor(lead); - var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount); + int incidentReminderCount = minIncidentNudgeCountFor(lead); + var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount, + incidentReminderCount); switch (action) { case INJECT -> { - injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount); + injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount, + incidentReminderCount); 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); } } @@ -480,11 +562,17 @@ 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 repliesBefore, Set ticketsBefore, - Set questionsBefore) { + Set questionsBefore) { + stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, pendingIncidentKeysFor(lead)); + } + + private void stopOrRestart(String lead, Set repliesBefore, Set ticketsBefore, + Set questionsBefore, Set incidentsBefore) { 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)); 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 +583,19 @@ 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) { // 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 replyTargets = pendingReplyTargetsFor(lead); List tickets = pendingTicketsFor(lead); List questions = pendingQuestionsFor(lead); - if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) { + List incidents = pendingIncidentsFor(lead); + if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.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); try { agents.send(lead, nudge); log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; " @@ -515,6 +604,11 @@ 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()); + } + } } catch (RuntimeException e) { log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}", lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, @@ -526,12 +620,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); + bumpNudgeCounts(replyTargets, tickets, questions, incidents); } /** Record that every one of these items was just named in a sent (or attempted) nudge. */ private void bumpNudgeCounts(Set replyTargets, List tickets, - List questions) { + List questions, List incidents) { for (String target : replyTargets) { pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1)); } @@ -544,6 +638,10 @@ 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)); + } } /** Schedule the next tick on the scheduler thread pool. */ @@ -556,7 +654,7 @@ public final class ReplyPushLoop { /** Render everything pending for one lead as a single nudge line. */ private static String formatNudge(Set replyTargets, List tickets, - List questions) { + List questions, List incidents) { List parts = new ArrayList<>(); if (!replyTargets.isEmpty()) { parts.add(formatRepliesNudge(replyTargets)); @@ -567,6 +665,9 @@ public final class ReplyPushLoop { if (!questions.isEmpty()) { parts.add(formatQuestionsNudge(questions)); } + if (!incidents.isEmpty()) { + parts.addAll(incidents.stream().map(ReplyPushLoop::formatBackendIncidentNudge).toList()); + } return String.join(" | ", parts); } @@ -606,6 +707,11 @@ 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())); + } + // --- lifecycle ----------------------------------------------------------------------------- /** @@ -627,6 +733,8 @@ public final class ReplyPushLoop { pendingReplies.clear(); pendingTickets.clear(); pendingQuestions.clear(); + pendingIncidents.clear(); + deliveredIncidents.clear(); } /** @see #stop() */ diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java index f449665..1d3bf54 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java @@ -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,103 @@ 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 = rec.sentParams().getFirst().getValue().toString(); + assertTrue(nudge.contains("unknown") && nudge.contains("no credential matched"), 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 +1026,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 From 7840e9adf6b339c972f1db2357feb95a21974976 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:44:18 +0700 Subject: [PATCH 2/2] CB-201: make unmapped target notice truthful --- .../dev/ltms/fleet/msg/ReplyPushLoop.java | 100 +++++++++++++++--- .../dev/ltms/fleet/msg/ReplyPushLoopTest.java | 7 +- 2 files changed, 91 insertions(+), 16 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java index c85228f..653d9c9 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -73,6 +73,9 @@ public final class ReplyPushLoop { 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; @@ -102,6 +105,9 @@ public final class ReplyPushLoop { private final ConcurrentHashMap pendingIncidents = new ConcurrentHashMap<>(); /** Incident/lead pairs already delivered. They make repeated incident reports one-shot. */ private final Set deliveredIncidents = ConcurrentHashMap.newKeySet(); + private final ConcurrentHashMap pendingUnmappedTargets = + new ConcurrentHashMap<>(); + private final Set deliveredUnmappedTargets = ConcurrentHashMap.newKeySet(); /** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */ private final ConcurrentHashMap activeLeads = new ConcurrentHashMap<>(); @@ -197,6 +203,12 @@ public final class ReplyPushLoop { List 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. */ private List pendingQuestionsFor(String lead) { return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList(); @@ -217,6 +229,15 @@ public final class ReplyPushLoop { .collect(Collectors.toUnmodifiableSet()); } + private List pendingUnmappedTargetsFor(String lead) { + return pendingUnmappedTargets.values().stream().filter(i -> lead.equals(i.key().lead())).toList(); + } + + private Set 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 minimum nudge count among the reply targets currently pending for it (CB-598). @@ -265,6 +286,11 @@ public final class ReplyPushLoop { 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. @@ -308,11 +334,19 @@ public final class ReplyPushLoop { /** 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(); - if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork) { + boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty(); + if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) { log.debug("push: nothing pending for lead {}, stopping reminder", lead); return Action.STOP; } @@ -320,7 +354,8 @@ public final class ReplyPushLoop { boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders; boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders; boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders; - if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible) { + 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"); @@ -468,8 +503,10 @@ public final class ReplyPushLoop { log.warn("push: backend target {} could not map to a credential: {}", target, reason); return; } - onBackendIncident("unmapped:" + target + ":" + reason, List.of(target), "unknown", - List.of("unmapped backend error: " + reason), 0); + 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 ---------------------------------------------------------------------------- @@ -504,21 +541,24 @@ public final class ReplyPushLoop { Set ticketsBefore = pendingTicketIdsFor(lead); Set questionsBefore = pendingQuestionTurnIdsFor(lead); Set incidentsBefore = pendingIncidentKeysFor(lead); + Set 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); + incidentReminderCount, unmappedTargetReminderCount); switch (action) { case INJECT -> { injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount, - incidentReminderCount); + 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, incidentsBefore); + case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, incidentsBefore, + unmappedTargetsBefore); } } @@ -568,11 +608,19 @@ public final class ReplyPushLoop { private void stopOrRestart(String lead, Set repliesBefore, Set ticketsBefore, Set questionsBefore, Set incidentsBefore) { + stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore, incidentsBefore, + pendingUnmappedTargetKeysFor(lead)); + } + + private void stopOrRestart(String lead, Set repliesBefore, Set ticketsBefore, + Set questionsBefore, Set incidentsBefore, + Set 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)) - || pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i)); + || 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); @@ -583,7 +631,8 @@ 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 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. @@ -591,11 +640,13 @@ public final class ReplyPushLoop { List tickets = pendingTicketsFor(lead); List questions = pendingQuestionsFor(lead); List incidents = pendingIncidentsFor(lead); - if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()) { + List 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, incidents); + String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets); try { agents.send(lead, nudge); log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; " @@ -609,6 +660,11 @@ public final class ReplyPushLoop { 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, @@ -620,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, incidents); + 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 replyTargets, List tickets, - List questions, List incidents) { + List questions, List incidents, + List unmappedTargets) { for (String target : replyTargets) { pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1)); } @@ -642,6 +699,10 @@ public final class ReplyPushLoop { 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. */ @@ -654,7 +715,8 @@ public final class ReplyPushLoop { /** Render everything pending for one lead as a single nudge line. */ private static String formatNudge(Set replyTargets, List tickets, - List questions, List incidents) { + List questions, List incidents, + List unmappedTargets) { List parts = new ArrayList<>(); if (!replyTargets.isEmpty()) { parts.add(formatRepliesNudge(replyTargets)); @@ -668,6 +730,9 @@ public final class ReplyPushLoop { 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); } @@ -712,6 +777,11 @@ public final class ReplyPushLoop { 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 ----------------------------------------------------------------------------- /** @@ -735,6 +805,8 @@ public final class ReplyPushLoop { pendingQuestions.clear(); pendingIncidents.clear(); deliveredIncidents.clear(); + pendingUnmappedTargets.clear(); + deliveredUnmappedTargets.clear(); } /** @see #stop() */ diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java index 1d3bf54..4d97c49 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java @@ -611,8 +611,11 @@ class ReplyPushLoopTest { loop.onBackendTargetUnmapped(WORKER, "no credential matched"); assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS)); - String nudge = rec.sentParams().getFirst().getValue().toString(); - assertTrue(nudge.contains("unknown") && nudge.contains("no credential matched"), nudge); + 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