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");