|
|
|
@@ -66,8 +66,13 @@ public final class ReplyPushLoop {
|
|
|
|
|
private final long backoffMs;
|
|
|
|
|
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
|
|
|
|
|
|
|
|
|
|
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
|
|
|
|
|
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
|
|
|
|
|
/**
|
|
|
|
|
* Worker targets with a reply queued, keyed by target. Each entry carries its own nudge
|
|
|
|
|
* count (CB-598) rather than sharing one counter per lead per source: a target's count only
|
|
|
|
|
* ever reflects nudges that actually named that target, so a target that joins while the
|
|
|
|
|
* schedule is already deep into another target's reminders still reads as fresh.
|
|
|
|
|
*/
|
|
|
|
|
private final ConcurrentHashMap<String, ReplyEntry> pendingReplies = new ConcurrentHashMap<>();
|
|
|
|
|
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
|
|
|
|
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
|
|
|
|
|
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
|
|
|
|
@@ -116,10 +121,10 @@ public final class ReplyPushLoop {
|
|
|
|
|
Set<String> result = new HashSet<>();
|
|
|
|
|
for (var entry : pendingReplies.entrySet()) {
|
|
|
|
|
String target = entry.getKey();
|
|
|
|
|
String owningLead = entry.getValue();
|
|
|
|
|
if (!lead.equals(owningLead)) continue;
|
|
|
|
|
ReplyEntry owning = entry.getValue();
|
|
|
|
|
if (!lead.equals(owning.lead())) continue;
|
|
|
|
|
if (inbox.peek(target).isEmpty()) {
|
|
|
|
|
pendingReplies.remove(target, owningLead);
|
|
|
|
|
pendingReplies.remove(target, owning);
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
result.add(target);
|
|
|
|
@@ -127,8 +132,15 @@ public final class ReplyPushLoop {
|
|
|
|
|
return result;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
|
|
|
|
|
private record PendingTicket(String ticket, String lead, boolean failed) {
|
|
|
|
|
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
|
|
|
|
|
private record ReplyEntry(String lead, int nudgeCount) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* A ticket awaiting collection: which lead to nudge, whether it ended in failure, and how
|
|
|
|
|
* many nudges have named it so far (CB-598 — tracked per ticket, not per lead per source).
|
|
|
|
|
*/
|
|
|
|
|
private record PendingTicket(String ticket, String lead, boolean failed, int nudgeCount) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
|
|
|
|
@@ -142,6 +154,41 @@ public final class ReplyPushLoop {
|
|
|
|
|
.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).
|
|
|
|
|
*
|
|
|
|
|
* <p>Before this, the count passed to {@code decide} was a single counter carried forward
|
|
|
|
|
* across scheduled ticks ({@code scheduleNext(lead, count + 1, ...)}), incremented whenever
|
|
|
|
|
* the source had <em>any</em> pending work — not tied to which target that work was. A target
|
|
|
|
|
* that joined while an older target's count was already near the cap inherited that count on
|
|
|
|
|
* its very next tick, even though no nudge had ever named it. Taking the minimum over what is
|
|
|
|
|
* actually pending now means a fresh target (count 0) keeps the source eligible regardless of
|
|
|
|
|
* how many times an older, still-undrained target has already been nudged; that older target
|
|
|
|
|
* keeps riding along in the combined nudge text without spending any more of its own budget
|
|
|
|
|
* (see {@link #bumpNudgeCounts}). Returns 0 when nothing is pending — {@link #decide} never
|
|
|
|
|
* consults the count in that case, since {@code hasReplyWork} is false.
|
|
|
|
|
*/
|
|
|
|
|
private int minReplyNudgeCountFor(String lead) {
|
|
|
|
|
int min = Integer.MAX_VALUE;
|
|
|
|
|
for (String target : pendingReplyTargetsFor(lead)) {
|
|
|
|
|
ReplyEntry entry = pendingReplies.get(target);
|
|
|
|
|
if (entry != null) {
|
|
|
|
|
min = Math.min(min, entry.nudgeCount());
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return min == Integer.MAX_VALUE ? 0 : min;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** As {@link #minReplyNudgeCountFor}, for the ticket source. */
|
|
|
|
|
private int minTicketNudgeCountFor(String lead) {
|
|
|
|
|
int min = Integer.MAX_VALUE;
|
|
|
|
|
for (PendingTicket ticket : pendingTicketsFor(lead)) {
|
|
|
|
|
min = Math.min(min, ticket.nudgeCount());
|
|
|
|
|
}
|
|
|
|
|
return min == Integer.MAX_VALUE ? 0 : min;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
|
|
|
|
* tickets alike — and return what the loop should do.
|
|
|
|
@@ -155,9 +202,14 @@ public final class ReplyPushLoop {
|
|
|
|
|
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
|
|
|
|
|
* {@link Action#STOP}.
|
|
|
|
|
*
|
|
|
|
|
* <p><strong>CB-598: the counts are per-item, not per-tick.</strong> {@link #tick} no longer
|
|
|
|
|
* carries these counts forward across scheduled calls — it recomputes them fresh every tick via
|
|
|
|
|
* {@link #minReplyNudgeCountFor} / {@link #minTicketNudgeCountFor}, so this function itself did
|
|
|
|
|
* not need to change; only what its caller feeds it did.
|
|
|
|
|
*
|
|
|
|
|
* @param lead the lead terminal to nudge
|
|
|
|
|
* @param replyReminderCount how many nudges have covered pending reply work for this lead
|
|
|
|
|
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
|
|
|
|
|
* @param replyReminderCount the lowest nudge count among reply targets pending for this lead
|
|
|
|
|
* @param ticketReminderCount the lowest nudge count among tickets pending for this lead
|
|
|
|
|
* @return the action the caller should take
|
|
|
|
|
*/
|
|
|
|
|
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
|
|
|
@@ -204,7 +256,8 @@ public final class ReplyPushLoop {
|
|
|
|
|
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
pendingReplies.put(target, lead.get());
|
|
|
|
|
pendingReplies.compute(target, (t, existing) ->
|
|
|
|
|
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
|
|
|
|
|
startOrCoalesce(lead.get());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -232,7 +285,8 @@ public final class ReplyPushLoop {
|
|
|
|
|
ticket, target);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
|
|
|
|
|
pendingTickets.compute(ticket, (id, existing) ->
|
|
|
|
|
new PendingTicket(ticket, lead.get(), failed, existing == null ? 0 : existing.nudgeCount()));
|
|
|
|
|
startOrCoalesce(lead.get());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -255,28 +309,37 @@ public final class ReplyPushLoop {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
log.debug("push: starting reminder loop for lead {}", lead);
|
|
|
|
|
scheduleNext(lead, 0, 0);
|
|
|
|
|
scheduleNext(lead);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Execute one loop tick — called on the scheduler thread. */
|
|
|
|
|
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
|
|
|
|
|
/**
|
|
|
|
|
* Execute one loop tick — called on the scheduler thread (or directly by a test; package-private
|
|
|
|
|
* for the same reason as {@link #stopOrRestart}).
|
|
|
|
|
*
|
|
|
|
|
* <p><strong>CB-598.</strong> The reminder counts fed into {@link #decide} are recomputed fresh
|
|
|
|
|
* every tick from what is actually pending right now ({@link #minReplyNudgeCountFor} /
|
|
|
|
|
* {@link #minTicketNudgeCountFor}), rather than carried forward as running counters across
|
|
|
|
|
* scheduled calls. A counter carried forward has no memory of which item it was counting for:
|
|
|
|
|
* a target or ticket that joined mid-backoff — after the previous tick fired but before this one
|
|
|
|
|
* did — is already sitting in {@code repliesBefore} / {@code ticketsBefore} below by the time this
|
|
|
|
|
* tick takes its snapshot, indistinguishable at that point from backlog the cap is meant to
|
|
|
|
|
* silence. Recomputing from the per-item counts fixes that: a newly-joined item's own count is
|
|
|
|
|
* still 0, so it keeps its source eligible regardless of how depleted an older, still-undrained
|
|
|
|
|
* item's count is.
|
|
|
|
|
*/
|
|
|
|
|
void tick(String lead) {
|
|
|
|
|
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
|
|
|
|
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
|
|
|
|
int replyReminderCount = minReplyNudgeCountFor(lead);
|
|
|
|
|
int ticketReminderCount = minTicketNudgeCountFor(lead);
|
|
|
|
|
var action = decide(lead, replyReminderCount, ticketReminderCount);
|
|
|
|
|
switch (action) {
|
|
|
|
|
case INJECT -> {
|
|
|
|
|
injectNudge(lead, replyReminderCount, ticketReminderCount);
|
|
|
|
|
// Only the source(s) actually eligible this tick spend a unit of their own budget —
|
|
|
|
|
// an exhausted source riding along in the combined message (still pending, still
|
|
|
|
|
// named) does not get charged again; its count stays put until it drains.
|
|
|
|
|
boolean replyEligible = !repliesBefore.isEmpty() && replyReminderCount < maxReminders;
|
|
|
|
|
boolean ticketEligible = !ticketsBefore.isEmpty() && ticketReminderCount < maxReminders;
|
|
|
|
|
scheduleNext(lead,
|
|
|
|
|
replyEligible ? replyReminderCount + 1 : replyReminderCount,
|
|
|
|
|
ticketEligible ? ticketReminderCount + 1 : ticketReminderCount);
|
|
|
|
|
scheduleNext(lead);
|
|
|
|
|
}
|
|
|
|
|
// Re-check after the configured backoff; the lead may become injectable soon.
|
|
|
|
|
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
|
|
|
|
|
case WAIT_BUSY -> scheduleNext(lead);
|
|
|
|
|
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@@ -317,7 +380,7 @@ public final class ReplyPushLoop {
|
|
|
|
|
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.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, 0, 0);
|
|
|
|
|
scheduleNext(lead);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
log.debug("push: reminder loop ended for lead {}", lead);
|
|
|
|
@@ -344,11 +407,28 @@ public final class ReplyPushLoop {
|
|
|
|
|
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
|
|
|
|
|
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
|
|
|
|
|
}
|
|
|
|
|
// Bump every item actually named in this nudge, not just whatever the shared source-level
|
|
|
|
|
// eligibility used to gate (CB-598) — each item's own count is what the next tick's
|
|
|
|
|
// minReplyNudgeCountFor / minTicketNudgeCountFor will read. An 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);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** 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) {
|
|
|
|
|
for (String target : replyTargets) {
|
|
|
|
|
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
|
|
|
|
}
|
|
|
|
|
for (PendingTicket ticket : tickets) {
|
|
|
|
|
pendingTickets.computeIfPresent(ticket.ticket(),
|
|
|
|
|
(id, e) -> new PendingTicket(e.ticket(), e.lead(), e.failed(), e.nudgeCount() + 1));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Schedule the next tick on the scheduler thread pool. */
|
|
|
|
|
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
|
|
|
|
|
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
|
|
|
|
|
private void scheduleNext(String lead) {
|
|
|
|
|
scheduler.schedule(() -> tick(lead),
|
|
|
|
|
backoffMs, TimeUnit.MILLISECONDS);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|