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