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 77aefb5..6e5bcdf 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -8,6 +8,8 @@ import dev.ltms.bridged.metrics.Metrics; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; @@ -16,35 +18,43 @@ import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; /** - * Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own - * herdr pane when a worker reply lands with no live {@code bridge_send} to resolve it. + * A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work + * waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an + * async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase + * (CB-588). * - *
The loop is triggered by {@link #onReplyQueued(String)} (called from - * {@link MessageService#reply} after the durable inbox publish). It checks four conditions - * at each tick via {@link #decide(String, int)}, then either injects a drain nudge, - * waits for the primary to become injectable, or stops reminding. + *
CB-590: one schedule per lead. Both kinds of work are triggered through + * their own entry point — {@link #onReplyQueued(String)} and + * {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be + * nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}. + * Earlier this was two independent schedules (one keyed by worker target for replies, one keyed + * by lead for tickets) that could both decide to inject into the same pane in the same window — + * a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn + * twice. Collapsing to one schedule per lead makes that structurally impossible: at most one + * scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}' + * compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight. * - *
Bounded: at most {@link #maxReminders} nudges per target, with a configurable backoff - * between them. The reply is never lost — the durable inbox is the backstop. - * - *
CB-588 ticket nudges. {@link #onTicketTerminal(String, String, boolean)} is a - * second, independent entry point for an async delegation ticket ({@code bridge_send(wait:false)}) - * reaching a terminal phase. That path takes the rendezvous fast path in {@link MessageService#reply} - * and never reaches {@link #onReplyQueued}, so without this the lead's own charter — prefer - * {@code wait:false} for anything non-trivial — was exactly the mode this loop failed to cover. It - * reuses the same status gating, bounded/backoff reminders, and metrics, keyed by the nudge-receiving - * lead terminal rather than the worker target so several tickets finishing together coalesce into one - * nudge. The two entry points do not interact: {@link #onReplyQueued} / {@link #decide} / their nudge - * text and bound are unchanged. + *
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, 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 {
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
- /** CB-588: singular form, one uncollected ticket. */
+ /** Coalesced form, several uncollected replies for the same lead. */
+ static final String REPLIES_NUDGE_FORMAT =
+ "%d workers returned replies — run bridge_poll(target=...) for each to collect them: %s";
+ /** Singular form, one uncollected ticket. */
static final String TICKET_NUDGE_FORMAT =
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
- /** CB-588: coalesced form, several uncollected tickets for the same lead. */
+ /** Coalesced form, several uncollected tickets for the same lead. */
static final String TICKETS_NUDGE_FORMAT =
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
@@ -56,11 +66,11 @@ public final class ReplyPushLoop {
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
- /** Track targets that have an active schedule. */
- private final ConcurrentHashMap Idempotent per lead: several tickets going terminal for the same lead while its schedule is
- * already active coalesce onto that schedule's next tick rather than firing a nudge each.
- *
- * @param ticket the ticket to nudge about
- * @param target the worker session the ticket was sent to — resolves which lead delegated it
- * @param failed whether the ticket ended in a failure phase rather than {@code DONE}
- */
- public void onTicketTerminal(String ticket, String target, boolean failed) {
- var lead = primaryRegistry.nudgeTargetFor(target);
- if (lead.isEmpty()) {
- log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
- ticket, target);
- return;
- }
- pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
- if (activeLeads.putIfAbsent(lead.get(), Boolean.TRUE) != null) {
- log.debug("push: ticket reminder loop already active for lead {}, {} coalesced in",
- lead.get(), ticket);
- return;
- }
- log.debug("push: starting ticket reminder loop for lead {}", lead.get());
- scheduleTicketTick(lead.get(), 0);
- }
-
- /**
- * Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
- * from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
- * lead already has (CB-588 acceptance #5). A ticket that was never pending (unknown ticket, or
- * one nudged with no push loop configured) is a no-op.
- */
- public void ticketCollected(String ticket) {
- pendingTickets.remove(ticket);
- }
-
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
- private List 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 decideTickets(String lead, int reminderCount) {
- if (pendingFor(lead).isEmpty()) {
- log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
+ 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: ticket 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;
}
@@ -278,95 +189,194 @@ public final class ReplyPushLoop {
return Action.WAIT_BUSY;
}
- /** Execute one ticket-loop tick — called on the scheduler thread. */
- private void ticketTick(String lead, int reminderCount) {
- Set Restarting on ANY non-empty {@code pendingFor(lead)} would be wrong: when STOP is reached
- * because the reminder cap was hit rather than the backlog draining, the same never-collected
- * ticket is expected to still be sitting there — that is the cap doing its job — and restarting
- * would nudge about it forever, defeating the bound (the original CB-307 bounded-reminder
- * guarantee, carried into CB-588 by acceptance criterion #7 — this exact regression showed up as
- * two existing tests failing once a naive "any pending ticket restarts" version of this fix went
- * in: {@code successfulTicketNudgeIncrementsDelivered} and {@code ticketNudgesSendUpToCapThenStop}).
- * Diffing the current pending set against {@code pendingBefore} tells the two cases apart: a
- * ticket present before this tick's decision is stale backlog, not a race; only a ticket absent
- * from {@code pendingBefore} can only have arrived during the decision-to-release window, which is
- * exactly the race this method closes.
+ * Resolves the delegating lead the same way {@link #onReplyQueued} does and coalesces onto
+ * the same per-lead schedule (CB-590) — several tickets, or a ticket and a reply, finishing
+ * for the same lead while its schedule is already active all ride the existing schedule's next
+ * tick rather than firing a nudge each.
+ *
+ * @param ticket the ticket to nudge about
+ * @param target the worker session the ticket was sent to — resolves which lead delegated it
+ * @param failed whether the ticket ended in a failure phase rather than {@code DONE}
+ */
+ public void onTicketTerminal(String ticket, String target, boolean failed) {
+ var lead = primaryRegistry.nudgeTargetFor(target);
+ if (lead.isEmpty()) {
+ log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
+ ticket, target);
+ return;
+ }
+ pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
+ startOrCoalesce(lead.get());
+ }
+
+ /**
+ * Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
+ * from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
+ * lead already has. A ticket that was never pending (unknown ticket, or one nudged with no push
+ * loop configured) is a no-op.
+ */
+ public void ticketCollected(String ticket) {
+ pendingTickets.remove(ticket);
+ }
+
+ // --- the schedule ----------------------------------------------------------------------------
+
+ /** Start a reminder schedule for {@code lead}, or join the one already running. */
+ private void startOrCoalesce(String lead) {
+ if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
+ log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
+ return;
+ }
+ log.debug("push: starting reminder loop for lead {}", lead);
+ scheduleNext(lead, 0, 0);
+ }
+
+ /** Execute one loop tick — called on the scheduler thread. */
+ private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
+ Set Restarting on ANY non-empty pending set would be wrong: when STOP is reached because the
+ * reminder cap was hit rather than the backlog draining, the same never-collected work is
+ * expected to still be sitting there — that is the cap doing its job — and restarting would
+ * nudge about it forever, defeating the bound. Diffing the current pending sets against the
+ * "before" snapshots tells the two cases apart: an item present before this tick's decision is
+ * stale backlog, not a race; only an item absent from the "before" snapshot can only have
+ * arrived during the decision-to-release window, which is exactly the race this method closes.
*
* Package-private so a test can drive the interleaving directly rather than trying to force a
- * genuine thread race: pass the exact {@code pendingBefore} snapshot a race requires (or does
- * not) and call this to prove the recheck responds correctly either way.
+ * genuine thread race.
*
- * Terminates rather than spinning: this method restarts the schedule at most once per call, and
- * a fresh {@link #onTicketTerminal} racing the recheck below still terminates in one of two ways —
- * either it observes the slot already vacated (by the {@code activeLeads.remove} above, which
- * happens-before this recheck in program order) and claims it itself, or it lands first and this
- * recheck then observes its ticket in {@code pendingTickets} and reclaims the slot instead. Exactly
- * one side always wins; neither can miss the other, so this never loops on its own account.
+ * Terminates rather than spinning: this method restarts the schedule at most once per call,
+ * and a fresh {@link #onReplyQueued} / {@link #onTicketTerminal} racing the recheck below still
+ * terminates in one of two ways — either it observes the slot already vacated (by the
+ * {@code activeLeads.remove} above, which happens-before this recheck in program order) and
+ * claims it itself, or it lands first and this recheck then observes its work in
+ * {@link #pendingReplies} / {@link #pendingTickets} and reclaims the slot instead. Exactly one
+ * side always wins; neither can miss the other, so this never loops on its own account.
*/
- void stopOrRestartTicketLoop(String lead, Set CB-590 collapsed the CB-307 reply-nudge schedule and the CB-588 ticket-nudge schedule into
+ * one schedule per lead ({@link ReplyPushLoop#decide}), so most tests below register pending work
+ * through the public entry points ({@code onReplyQueued} / {@code onTicketTerminal}) before
+ * exercising {@code decide} directly, mirroring how the two entry points now share one decision
+ * function keyed by the lead terminal rather than by worker target.
+ *
* Uses a {@link RecordingHerdrClient} that synchronizes access to its call list so the
* scheduler thread and test thread never have memory ordering issues. The {@code decide()}
* tests use a simple client with no concurrency concern.
@@ -58,38 +65,50 @@ class ReplyPushLoopTest {
// --- decide() logic ------------------------------------------------------------------------
@Test
- void decideWithoutPrimaryIsStop() {
- agents = agentWithStatus("idle");
- var loop = new ReplyPushLoop(
- new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100);
- assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0));
+ void onReplyQueuedWithNoKnownLeadNeverStartsASchedule() throws Exception {
+ var rec = recordingClient();
+ agents = new AgentControl(rec);
+ inbox.publish(WORKER, "m1", "hello");
+ var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
+
+ loop.onReplyQueued(WORKER); // no lead known -> never registered, never scheduled
+
+ assertFalse(loop.isActive(), "no lead known means nothing to nudge yet");
+ Thread.sleep(150);
+ assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting");
}
@Test
- void decideWithEmptyInboxIsStop() {
+ void decideWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
- assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
+ assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
void decideAtCapIsStop() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
+ var loop = loop(2, 100);
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
}
@Test
void decideUnderCapWithInjectablePrimaryIsInject() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
+ var loop = loop();
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideUnderCapWithBlockedPrimaryIsInject() {
agents = agentWithStatus("blocked");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
+ var loop = loop();
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"BLOCKED is injectable");
}
@@ -97,7 +116,9 @@ class ReplyPushLoopTest {
void decideUnderCapWithDonePrimaryIsInject() {
agents = agentWithStatus("done");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
+ var loop = loop();
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"DONE is injectable");
}
@@ -105,23 +126,29 @@ class ReplyPushLoopTest {
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
agents = agentWithStatus("working");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
+ var loop = loop();
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
agents = agentWithStatus("unknown");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
+ var loop = loop();
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideStopsAfterInboxIsEmptied() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
- assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
+ var loop = loop();
+ loop.onReplyQueued(WORKER);
+ assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
inbox.ack(WORKER, "m1");
- assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
+ assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@@ -201,12 +228,19 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
}
- // --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
+ @Test
+ void repliesNudgeFormatIsCorrect() {
+ String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
+ assertTrue(multi.contains("2 workers"));
+ assertTrue(multi.contains("bridge_poll(target=...)"));
+ }
+
+ // --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
@Test
void decideTicketsWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
- assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
+ assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
@@ -214,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.decideTickets(PRIMARY, 2));
+ assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
}
@Test
@@ -222,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.decideTickets(PRIMARY, 0));
+ assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -230,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.decideTickets(PRIMARY, 0));
+ assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -238,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.decideTickets(PRIMARY, 0));
+ assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
@@ -255,7 +289,7 @@ class ReplyPushLoopTest {
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call");
- assertFalse(nudge.contains("bridge_poll(target="), "a ticket nudge must not tell the lead to run the target-poll call");
+ assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
}
@Test
@@ -346,11 +380,11 @@ class ReplyPushLoopTest {
void aTicketStillPendingWhenTheLoopStopsIsNotStranded() {
// Regression for the race a reviewer found in gitea PR #73: onTicketTerminal's
// activeLeads.putIfAbsent can see the lead's slot as still occupied a moment before
- // decideTickets' STOP releases it, so the ticket coalesces onto a schedule that is about to
+ // decide's STOP releases it, so the ticket coalesces onto a schedule that is about to
// die and nothing ever nudges about it. Forcing that exact thread interleaving is not
- // reliable, so this drives stopOrRestartTicketLoop — the STOP path's own release-and-recheck —
+ // reliable, so this drives stopOrRestart — the STOP path's own release-and-recheck —
// directly, arranging the state it must not lose a ticket in: a ticket pending for the lead
- // that was NOT part of the pre-decision snapshot (pendingBefore=empty), standing in for one
+ // that was NOT part of the pre-decision snapshot (ticketsBefore=empty), standing in for one
// that races in during the decision-to-release window.
var rec = recordingClient();
agents = new AgentControl(rec);
@@ -358,9 +392,9 @@ class ReplyPushLoopTest {
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
- // Stand in for the scheduler thread reaching decideTickets==STOP for this lead — with nothing
+ // Stand in for the scheduler thread reaching decide==STOP for this lead — with nothing
// pending at decide time — while task-1 races in before the release below runs.
- loop.stopOrRestartTicketLoop(PRIMARY, Set.of());
+ loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
assertTrue(loop.isActive(), "a ticket that raced the loop's stop must reclaim the schedule "
+ "slot, not be stranded with no schedule left to ever nudge about it");
@@ -368,23 +402,67 @@ class ReplyPushLoopTest {
@Test
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
- // The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) would be
+ // The other direction of the same fix: restarting on ANY non-empty pending set would be
// wrong. When STOP is reached because the reminder cap was hit, the same never-collected
// ticket is expected to still be there — that is the cap doing its job (acceptance criterion
- // #7: nudges stay bounded). task-1 here was already accounted for at decide time (it is in
- // pendingBefore), so it must not restart the loop just because it is still sitting there.
+ // #5: no spin / nudges stay bounded). task-1 here was already accounted for at decide time
+ // (it is in ticketsBefore), so it must not restart the loop just because it is still sitting
+ // there.
var rec = recordingClient();
agents = new AgentControl(rec);
ReplyPushLoop loop = loop(1, 100_000);
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
- loop.stopOrRestartTicketLoop(PRIMARY, Set.of("task-1"));
+ loop.stopOrRestart(PRIMARY, Set.of(), Set.of("task-1"));
assertFalse(loop.isActive(), "a stale ticket already accounted for at decide time must not "
+ "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");
@@ -402,7 +480,103 @@ class ReplyPushLoopTest {
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
- "CB-588 must not change the CB-307 inbox nudge's call shape");
+ "CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
+ }
+
+ // --- CB-590: one schedule per lead — no overlap, no lost nudges -----------------------------
+
+ @Test
+ void replyAndTicketForTheSameLeadCoalesceIntoOneSendNeverOverlapping() throws Exception {
+ var rec = recordingClient();
+ agents = new AgentControl(rec);
+ inbox.publish(WORKER, "m1", "hello");
+ var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
+
+ loop.onReplyQueued(WORKER);
+ loop.onTicketTerminal("task-1", WORKER, false);
+
+ assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
+ Thread.sleep(300);
+ assertEquals(1, rec.sendCount(),
+ "a reply and a ticket for the same lead must coalesce onto ONE schedule — "
+ + "two nudge injections into the same lead pane must never overlap");
+ String nudge = rec.sentParams().getFirst().getValue().toString();
+ assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
+ "the combined nudge must still mention the reply: " + nudge);
+ assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
+ }
+
+ @Test
+ void aReplyQueuedWhileTheLeadIsBusyIsNotLostWhenATicketArrivesToo() throws Exception {
+ // The lead is busy for its first two status checks, then becomes injectable. A reply is
+ // queued while busy; a ticket for the same lead arrives before the lead frees up. Neither
+ // may be dropped — deferred is fine, lost is not (acceptance criterion #2).
+ var rec = new BusyThenIdleHerdrClient(2);
+ agents = new AgentControl(rec);
+ inbox.publish(WORKER, "m1", "hello");
+ var loop = loop(1, 50); // cap=1: WAIT_BUSY doesn't count against it, so exactly one send once injectable
+
+ loop.onReplyQueued(WORKER); // schedule starts, first tick(s) WAIT_BUSY
+ loop.onTicketTerminal("task-1", WORKER, false); // coalesces onto the same waiting schedule
+
+ assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
+ "once the lead becomes injectable, the deferred work must still be nudged");
+ Thread.sleep(200);
+ assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
+ String nudge = rec.sentParams().getFirst().getValue().toString();
+ assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
+ 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) ----------------------------------------------------------------------
@@ -430,8 +604,10 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
+ var loop = loop(2, 100, metrics);
+ loop.onReplyQueued(WORKER);
- assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 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");
@@ -459,7 +635,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100_000, metrics);
loop.onTicketTerminal("task-1", WORKER, false);
- assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(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");
@@ -547,4 +723,49 @@ class ReplyPushLoopTest {
private static RecordingHerdrClient recordingClient() {
return new RecordingHerdrClient();
}
+
+ /**
+ * Thread-safe fake that reports {@code working} (not injectable) for its first
+ * {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
+ * while the lead is busy is deferred, not dropped, once it becomes injectable.
+ */
+ private static final class BusyThenIdleHerdrClient implements HerdrClient {
+ private final List