From 7840e9adf6b339c972f1db2357feb95a21974976 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:44:18 +0700 Subject: [PATCH] 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