CB-590 follow-up: give each nudge source its own reminder budget
decide(lead, reminderCount) shared one counter across the reply and ticket sources after PR #84 collapsed both onto a single per-lead schedule. A reply stream that used up the whole budget could then make decide() STOP even for a ticket that had never been nudged and had coalesced onto the same still-active schedule — stranding it with no live schedule left, since stopOrRestart's racedIn check does not save work that was already present in the "before" snapshot. decide() now tracks a per-source count (replyReminderCount, ticketReminderCount) and returns INJECT while either source is still under its own cap, STOP only when both are exhausted. Still exactly one schedule per lead; stopOrRestart's snapshot-diff logic is untouched.
This commit is contained in:
@@ -37,10 +37,12 @@ import java.util.stream.Collectors;
|
||||
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
|
||||
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
|
||||
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
|
||||
* ({@link #injectNudge(String, int)}). Work that arrives while the lead is busy is never lost:
|
||||
* it is re-read fresh on every tick until the lead is injectable or the shared reminder cap
|
||||
* ({@link #maxReminders}) is reached, whichever the durable inbox / pending-ticket set doesn't
|
||||
* already answer via {@code STOP}.
|
||||
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
|
||||
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
|
||||
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
|
||||
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
|
||||
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
|
||||
* via {@code STOP}.
|
||||
*/
|
||||
public final class ReplyPushLoop {
|
||||
|
||||
@@ -144,20 +146,32 @@ public final class ReplyPushLoop {
|
||||
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
||||
* tickets alike — and return what the loop should do.
|
||||
*
|
||||
* @param lead the lead terminal to nudge
|
||||
* @param reminderCount how many nudges have been sent so far for this lead (shared across
|
||||
* both reply and ticket work — CB-590 collapses the reminder cap onto
|
||||
* one counter per lead, so alternating sources cannot outrun the bound)
|
||||
* <p><strong>CB-590-fix: one schedule, two budgets.</strong> The single per-lead schedule
|
||||
* (CB-590) still ticks once for both sources, but each source is capped independently —
|
||||
* {@code replyReminderCount} against a reply target still pending, {@code ticketReminderCount}
|
||||
* against a ticket still pending. A busy reply stream that exhausts its own cap must not stop
|
||||
* the loop from nudging about a ticket that still has budget left, and vice versa: either
|
||||
* source being eligible (has pending work AND is under its own cap) is enough for
|
||||
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
|
||||
* {@link Action#STOP}.
|
||||
*
|
||||
* @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
|
||||
* @return the action the caller should take
|
||||
*/
|
||||
Action decide(String lead, int reminderCount) {
|
||||
boolean anyPending = !pendingReplyTargetsFor(lead).isEmpty() || !pendingTicketIdsFor(lead).isEmpty();
|
||||
if (!anyPending) {
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
||||
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork) {
|
||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||
return Action.STOP;
|
||||
}
|
||||
if (reminderCount >= maxReminders) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
|
||||
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
||||
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||
maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
return Action.STOP;
|
||||
}
|
||||
@@ -241,21 +255,28 @@ public final class ReplyPushLoop {
|
||||
return;
|
||||
}
|
||||
log.debug("push: starting reminder loop for lead {}", lead);
|
||||
scheduleNext(lead, 0);
|
||||
scheduleNext(lead, 0, 0);
|
||||
}
|
||||
|
||||
/** Execute one loop tick — called on the scheduler thread. */
|
||||
private void tick(String lead, int reminderCount) {
|
||||
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
||||
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
||||
var action = decide(lead, reminderCount);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(lead, reminderCount);
|
||||
scheduleNext(lead, reminderCount + 1);
|
||||
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);
|
||||
}
|
||||
// Re-check after the configured backoff; the lead may become injectable soon.
|
||||
case WAIT_BUSY -> scheduleNext(lead, reminderCount);
|
||||
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
||||
}
|
||||
}
|
||||
@@ -296,14 +317,14 @@ 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);
|
||||
scheduleNext(lead, 0, 0);
|
||||
return;
|
||||
}
|
||||
log.debug("push: reminder loop ended for lead {}", lead);
|
||||
}
|
||||
|
||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||
private void injectNudge(String lead, int reminderCount) {
|
||||
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
|
||||
// collected (or another arrive), between the decision and the injection.
|
||||
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
||||
@@ -315,18 +336,20 @@ public final class ReplyPushLoop {
|
||||
String nudge = formatNudge(replyTargets, tickets);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: nudge {}/{} sent to lead {} ({} reply target(s), {} ticket(s))",
|
||||
reminderCount + 1, maxReminders, lead, replyTargets.size(), tickets.size());
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
replyTargets.size(), tickets.size());
|
||||
countNudge("delivered");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} (reminder {}/{}): {}",
|
||||
lead, reminderCount + 1, maxReminders, e.toString());
|
||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
|
||||
}
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext(String lead, int nextReminderCount) {
|
||||
scheduler.schedule(() -> tick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
|
||||
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
|
||||
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
|
||||
backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
// --- nudge formatting ------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user