From 88b9503c3beab253c394784ca9d63bce3659a430 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 16 Aug 2026 17:26:07 +0200 Subject: [PATCH] CB-590 follow-up: give each nudge source its own reminder budget MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../dev/ltms/bridged/msg/ReplyPushLoop.java | 77 +++++++---- .../ltms/bridged/msg/ReplyPushLoopTest.java | 126 +++++++++++++++--- 2 files changed, 160 insertions(+), 43 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index f368f1e..6e5bcdf 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -37,10 +37,12 @@ import java.util.stream.Collectors; *

Each tick examines everything 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) + *

CB-590-fix: one schedule, two budgets. 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 repliesBefore = pendingReplyTargetsFor(lead); Set 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 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 ------------------------------------------------------------------------ diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java index 974245d..ad58a15 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -81,7 +81,7 @@ class ReplyPushLoopTest { @Test void decideWithNothingPendingIsStop() { agents = agentWithStatus("idle"); - assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0)); } @Test @@ -90,7 +90,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(2, 100); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0)); } @Test @@ -99,7 +99,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0)); } @Test @@ -108,7 +108,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0), + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0), "BLOCKED is injectable"); } @@ -118,7 +118,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0), + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0), "DONE is injectable"); } @@ -128,7 +128,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0)); } @Test @@ -137,7 +137,7 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0)); } @Test @@ -146,9 +146,9 @@ class ReplyPushLoopTest { inbox.publish(WORKER, "m1", "hello"); var loop = loop(); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0)); inbox.ack(WORKER, "m1"); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0)); } // --- onReplyQueued integration ------------------------------------------------------------- @@ -240,7 +240,7 @@ class ReplyPushLoopTest { @Test void decideTicketsWithNothingPendingIsStop() { agents = agentWithStatus("idle"); - assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0)); } @Test @@ -248,7 +248,7 @@ class ReplyPushLoopTest { agents = agentWithStatus("idle"); var loop = loop(2, 100_000); loop.onTicketTerminal("task-1", WORKER, false); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2)); } @Test @@ -256,7 +256,7 @@ class ReplyPushLoopTest { agents = agentWithStatus("idle"); var loop = loop(5, 100_000); loop.onTicketTerminal("task-1", WORKER, false); - assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0)); } @Test @@ -264,7 +264,7 @@ class ReplyPushLoopTest { agents = agentWithStatus("working"); var loop = loop(5, 100_000); loop.onTicketTerminal("task-1", WORKER, false); - assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0)); } @Test @@ -272,7 +272,7 @@ class ReplyPushLoopTest { agents = agentWithStatus("idle"); var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100_000); loop.onTicketTerminal("task-1", WORKER, false); // no lead known -> never registered as pending - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0)); } // --- CB-588: onTicketTerminal integration --------------------------------------------------- @@ -420,6 +420,49 @@ class ReplyPushLoopTest { + "restart the loop — that would defeat the reminder cap"); } + // --- mirror of the two stopOrRestart tests above, for the reply arm ------------------------ + // + // Both tests above only ever passed Set.of() for repliesBefore, so racedIn's reply branch + // (`pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))`) was + // never exercised by anything other than an always-empty snapshot. The reviewer flagged this: + // racedIn is symmetric in the code, and only half of it was pinned by a test. + + @Test + void aReplyStillPendingWhenTheLoopStopsIsNotStranded() { + // Mirrors aTicketStillPendingWhenTheLoopStopsIsNotStranded: a reply target that raced in + // during the decision-to-release window (absent from the "before" snapshot) must reclaim + // the schedule slot rather than being stranded with no schedule left to nudge about it. + var rec = recordingClient(); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + ReplyPushLoop loop = loop(1, 100_000); // long backoff — no natural tick fires during this test + + loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY} + + loop.stopOrRestart(PRIMARY, Set.of(), Set.of()); + + assertTrue(loop.isActive(), "a reply that raced the loop's stop must reclaim the schedule " + + "slot, not be stranded with no schedule left to ever nudge about it"); + } + + @Test + void aStaleUncollectedReplyAtCapDoesNotRestartTheLoop() { + // Mirrors aStaleUncollectedTicketAtCapDoesNotRestartTheLoop: a reply target already + // accounted for at decide time (present in repliesBefore) must not restart the loop — + // that is the reminder cap doing its job, not a race. + var rec = recordingClient(); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + ReplyPushLoop loop = loop(1, 100_000); + + loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY} + + loop.stopOrRestart(PRIMARY, Set.of(WORKER), Set.of()); + + assertFalse(loop.isActive(), "a stale reply target already accounted for at decide time " + + "must not restart the loop — that would defeat the reminder cap"); + } + @Test void ticketNudgeFormatIsCorrect() { String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1"); @@ -485,6 +528,57 @@ class ReplyPushLoopTest { assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge); } + // --- CB-590 follow-up: per-source reminder budgets — the regression this round exists for --- + + @Test + void oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource() { + // The live trace this ticket was filed from: an undrained reply target got nudged up to + // its cap (5 reminders), then a ticket for the SAME lead went terminal shortly before the + // next scheduled tick — so it coalesced onto the still-active schedule (arriving BEFORE + // that tick's "before" snapshot, not during the decision-to-release race stopOrRestart + // guards). With CB-590's single shared reminder counter, that tick's decide() saw + // reminderCount already at the cap and returned STOP regardless of the ticket, and because + // the ticket was already present in that tick's "before" snapshot, stopOrRestart's + // racedIn check (proven correct on its own above) did not save it either — it is not a + // race, it looks like ordinary stale backlog. The ticket was then stranded: pending + // forever with no live schedule, never named in any nudge. + // + // Fixed by giving each source its own counter. Here the reply source is AT its cap (2/2) + // and the ticket source has NEVER been nudged (0/2) — decide() must still return INJECT, + // because the ticket is still eligible on its own budget. + agents = agentWithStatus("idle"); + inbox.publish(WORKER, "m1", "hello"); + var loop = loop(2, 100_000); // huge backoff — this test drives decide()/isActive() directly + loop.onReplyQueued(WORKER); + loop.onTicketTerminal("task-1", WORKER, false); + + assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 2, 0), + "the reply source is exhausted (2/2), but the ticket source has never been " + + "nudged (0/2) — the lead must still be injected so the ticket is not " + + "lost, exactly the CB-590 follow-up regression"); + + // isActive() (criterion #5): a real tick that takes the INJECT branch above never calls + // stopOrRestart, so the schedule started by onTicketTerminal above stays live — the + // ticket is not left stranded with isActive()==false while it is still pending. + assertTrue(loop.isActive(), "the schedule must stay active while the ticket source still " + + "has budget left, even though the reply source sharing it is exhausted"); + } + + @Test + void bothSourcesExhaustedIsStillStop() { + // The flip side: per-source budgets must not turn into unbounded nudging. When BOTH + // sources are at their cap, decide() must still STOP — a per-source budget is still a + // budget. + agents = agentWithStatus("idle"); + inbox.publish(WORKER, "m1", "hello"); + var loop = loop(2, 100_000); + loop.onReplyQueued(WORKER); + loop.onTicketTerminal("task-1", WORKER, false); + + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 2), + "both the reply and the ticket source are at their own cap — must still stop"); + } + // --- metrics (CB-512) ---------------------------------------------------------------------- @Test @@ -513,7 +607,7 @@ class ReplyPushLoopTest { var loop = loop(2, 100, metrics); loop.onReplyQueued(WORKER); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0)); assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"), "hitting the reminder cap must count as exhausted"); @@ -541,7 +635,7 @@ class ReplyPushLoopTest { var loop = loop(2, 100_000, metrics); loop.onTicketTerminal("task-1", WORKER, false); - assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2)); + assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2)); assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"), "hitting the ticket reminder cap must count as exhausted");