From 81a0cf471022e53bc662b0bad6e454d863b9413a Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 16 Aug 2026 17:49:21 +0200 Subject: [PATCH] CB-598: track reminder counts per pending item, not per lead per source Work that arrived during the ~15s push_backoff_ms window between two ticks landed in the pending map before the next tick's start-of-tick snapshot, so a shared per-lead-per-source counter (carried forward via scheduleNext(lead, count+1, ...)) already treated it as exhausted backlog even though no nudge had ever named it. ReplyPushLoop.tick now recomputes each source's reminder count fresh every tick as the minimum nudge count among that source's currently pending items, so a freshly-arrived item (count 0) keeps its source eligible regardless of how depleted an older, still-undrained sibling's count is. decide() itself is unchanged. --- .../dev/ltms/bridged/msg/ReplyPushLoop.java | 132 ++++++++++++++---- .../ltms/bridged/msg/ReplyPushLoopTest.java | 78 +++++++++++ 2 files changed, 184 insertions(+), 26 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 6e5bcdf..164b618 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -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 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 pendingReplies = new ConcurrentHashMap<>(); /** Tickets that have gone terminal but not yet been polled, keyed by ticket. */ private final ConcurrentHashMap 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 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 minimum nudge count among the reply targets currently pending for it (CB-598). + * + *

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 any 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}. * + *

CB-598: the counts are per-item, not per-tick. {@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}). + * + *

CB-598. 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 repliesBefore = pendingReplyTargetsFor(lead); Set 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 replyTargets, List 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); } 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 ad58a15..71557ab 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -42,6 +42,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 ObjectMapper MAPPER = new ObjectMapper(); private PrimaryRegistry registry; @@ -579,6 +580,83 @@ class ReplyPushLoopTest { "both the reply and the ticket source are at their own cap — must still stop"); } + // --- CB-598: work arriving during a backoff must not read as stale backlog ------------------- + + @Test + void aTargetArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() { + // The bug: reminder counts used to be a single counter per lead per source, carried + // forward across scheduled ticks (scheduleNext(lead, count + 1, ...)) rather than tracked + // per pending item. WORKER gets nudged once here, which — with cap=1 — exhausts the + // shared reply-source counter for this lead. WORKER2 then queues a reply for the SAME + // lead "during the backoff": while the schedule from WORKER's tick is still active, before + // the next tick's own start-of-tick snapshot runs. At that next tick, the OLD code passed + // the already-exhausted shared counter into decide() regardless of WORKER2 never having + // been named in any nudge, and — because WORKER2 was already present in that tick's + // "before" snapshot — stopOrRestart's race check (proven correct on its own elsewhere in + // this file) does not save it either: it looks like ordinary stale backlog, not a race. + // WORKER2 was then stranded forever with no live schedule and no nudge ever naming it. + // + // tick() is driven directly (package-private, same reasoning as stopOrRestart being + // directly testable) so the exact interleaving is deterministic instead of racing the + // scheduler thread over a real ~15s backoff. + // + // Before the fix, this test fails on the second assertEquals: rec.sendCount() stays at 1 + // (decide() returns STOP on the second tick(), so injectNudge is never called a second + // time) and the "must still get one" assertion never even runs. + int cap = 1; + var rec = recordingClient(); + agents = new AgentControl(rec); + inbox.own(WORKER2); + inbox.publish(WORKER, "m1", "hello"); + var loop = loop(cap, 100_000); // huge backoff — nothing fires on its own; we drive tick() + + loop.onReplyQueued(WORKER); + loop.tick(PRIMARY); // first tick: nudges WORKER alone; WORKER's own count reaches the cap + assertEquals(1, rec.sendCount(), "the first tick should nudge about WORKER"); + + // WORKER2 "arrives during the backoff": queued for the same lead while the schedule from + // the tick above is still active (activeLeads still holds PRIMARY), before the next tick + // (simulated below) takes its own start-of-tick snapshot. + inbox.publish(WORKER2, "m2", "hello2"); + loop.onReplyQueued(WORKER2); + + loop.tick(PRIMARY); // the tick that would fire once that backoff elapsed + + assertEquals(2, rec.sendCount(), + "WORKER2 was never named in any nudge yet and must still get one, even though " + + "WORKER's own reminder count is already at the cap"); + String secondNudge = rec.sentParams().get(1).getValue().toString(); + assertTrue(secondNudge.contains(WORKER2), "the never-named target must be named: " + secondNudge); + + // Criterion #3: isActive() must reflect that this lead still had a live nudge to give — + // the second tick took the INJECT branch, so the schedule stayed live rather than being + // torn down under WORKER2. + assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh target"); + } + + @Test + void aTicketArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() { + // Mirrors the reply-side test above for the ticket source. + int cap = 1; + var rec = recordingClient(); + agents = new AgentControl(rec); + var loop = loop(cap, 100_000); + + loop.onTicketTerminal("task-1", WORKER, false); + loop.tick(PRIMARY); // first tick: nudges task-1 alone; its count reaches the cap + assertEquals(1, rec.sendCount(), "the first tick should nudge about task-1"); + + loop.onTicketTerminal("task-2", WORKER, false); // arrives during the backoff, same lead + loop.tick(PRIMARY); + + assertEquals(2, rec.sendCount(), + "task-2 was never named in any nudge yet and must still get one, even though " + + "task-1's reminder count is already at the cap"); + String secondNudge = rec.sentParams().get(1).getValue().toString(); + assertTrue(secondNudge.contains("task-2"), "the never-named ticket must be named: " + secondNudge); + assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket"); + } + // --- metrics (CB-512) ---------------------------------------------------------------------- @Test