Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 16de9df000 CB-601: make the recovery-race test's head start deterministic, not a sleep
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Successful in 1m9s
2026-08-16 18:02:57 +02:00
3 changed files with 63 additions and 197 deletions
@@ -66,13 +66,8 @@ 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, 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<String, ReplyEntry> pendingReplies = new ConcurrentHashMap<>();
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
@@ -121,10 +116,10 @@ public final class ReplyPushLoop {
Set<String> result = new HashSet<>();
for (var entry : pendingReplies.entrySet()) {
String target = entry.getKey();
ReplyEntry owning = entry.getValue();
if (!lead.equals(owning.lead())) continue;
String owningLead = entry.getValue();
if (!lead.equals(owningLead)) continue;
if (inbox.peek(target).isEmpty()) {
pendingReplies.remove(target, owning);
pendingReplies.remove(target, owningLead);
continue;
}
result.add(target);
@@ -132,15 +127,8 @@ public final class ReplyPushLoop {
return result;
}
/** 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) {
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
private record PendingTicket(String ticket, String lead, boolean failed) {
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
@@ -154,41 +142,6 @@ public final class ReplyPushLoop {
.collect(Collectors.toUnmodifiableSet());
}
/**
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
*
* <p>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 <em>any</em> 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.
@@ -202,14 +155,9 @@ public final class ReplyPushLoop {
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
* {@link Action#STOP}.
*
* <p><strong>CB-598: the counts are per-item, not per-tick.</strong> {@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 the lowest nudge count among reply targets pending for this lead
* @param ticketReminderCount the lowest nudge count among tickets pending for this lead
* @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 replyReminderCount, int ticketReminderCount) {
@@ -256,8 +204,7 @@ public final class ReplyPushLoop {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
return;
}
pendingReplies.compute(target, (t, existing) ->
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
pendingReplies.put(target, lead.get());
startOrCoalesce(lead.get());
}
@@ -285,8 +232,7 @@ public final class ReplyPushLoop {
ticket, target);
return;
}
pendingTickets.compute(ticket, (id, existing) ->
new PendingTicket(ticket, lead.get(), failed, existing == null ? 0 : existing.nudgeCount()));
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
startOrCoalesce(lead.get());
}
@@ -309,37 +255,28 @@ public final class ReplyPushLoop {
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead);
scheduleNext(lead, 0, 0);
}
/**
* Execute one loop tick — called on the scheduler thread (or directly by a test; package-private
* for the same reason as {@link #stopOrRestart}).
*
* <p><strong>CB-598.</strong> 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) {
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
Set<String> 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);
scheduleNext(lead);
// 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);
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
}
}
@@ -380,7 +317,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);
scheduleNext(lead, 0, 0);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
@@ -407,28 +344,11 @@ 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<String> replyTargets, List<PendingTicket> 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) {
scheduler.schedule(() -> tick(lead),
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
backoffMs, TimeUnit.MILLISECONDS);
}
@@ -35,9 +35,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
*/
class AmqpReplyInboxRecoveryRaceTest {
/** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
private static final int STALE_PUBLISHES = 100_000;
/** Large enough that thousands of entries are still unprocessed by the time the very first one
* is observed as failed (see {@code sweepIsHoldingTheLock} below) — that gap is what makes the
* head start deterministic instead of a coin flip. 100,000 gave the same guarantee but made the
* test far more expensive than the guarantee needs; the ordering no longer depends on a timing
* window sized to the full backlog; just to the tail of it. */
private static final int STALE_PUBLISHES = 2_000;
@Test
@Timeout(30)
@@ -58,6 +61,12 @@ class AmqpReplyInboxRecoveryRaceTest {
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
// Virtual threads make this many concurrent blocking publish() calls cheap.
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
// Counted down by the FIRST stale publish thread to observe its own failure. That can only
// happen from inside failPendingPublishesOnRecovery() — nothing else in this test ever
// completes a stale Pending exceptionally (no nack/return is simulated for any "stale-*"
// msgId) — so seeing it fire is direct, observable proof the sweep is inside its loop, not a
// timing guess. It replaces the old fixed Thread.sleep(5) head start.
CountDownLatch sweepIsHoldingTheLock = new CountDownLatch(1);
for (int i = 0; i < STALE_PUBLISHES; i++) {
String msgId = "stale-" + i;
Thread.ofVirtual().start(() -> {
@@ -65,7 +74,7 @@ class AmqpReplyInboxRecoveryRaceTest {
try {
inbox.publish("worker-stale", msgId, "x");
} catch (IllegalStateException expected) {
// resolved (failed by the sweep) — that is exactly what this thread is here for
sweepIsHoldingTheLock.countDown();
}
});
}
@@ -83,13 +92,28 @@ class AmqpReplyInboxRecoveryRaceTest {
}
}, "recovery-sweep");
sweepThread.start();
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
// register. Without this head start, "fresh" sometimes wins the race for the lock and
// registers before the sweep even starts, which is the accepted "already in flight when
// recovery fires" case (correctly failed either way) rather than the bug under test.
Thread.sleep(5);
// Deterministic head start: block until the sweep has actually failed one of the stale
// publishes. failPendingPublishesOnRecovery() (once guarded, as it is on main) holds
// publishChannelLock for its ENTIRE loop, not just per entry — so this failure proves the
// sweep is, at this instant, still holding that lock. With STALE_PUBLISHES this large, the
// remaining ~1,999 entries give an enormous margin between "first failure observed" and "sweep
// releases the lock": there is no window left for "fresh" to slip in before the sweep starts,
// or to win the lock ahead of it — see the case-2 note below. This also means Case 1 (the sweep
// is already inside its loop, holding the lock, when "fresh" tries to register) is now
// guaranteed by construction rather than merely likely under a fixed sleep.
assertTrue(sweepIsHoldingTheLock.await(20, TimeUnit.SECONDS),
"the sweep never failed a single stale publish — it may not have started");
// Case 2 ("fresh" wins publishChannelLock before the sweep even starts, so it genuinely
// published on the stale channel and the sweep correctly fails it) is impossible by
// construction in this test: freshThread.start() below is reached only after
// sweepIsHoldingTheLock has counted down, which can only happen once
// failPendingPublishesOnRecovery() is already running and has already failed a stale entry.
// There is no code path that lets "fresh" start before the sweep starts. That case is real
// and correct production behaviour (see AmqpReplyInbox#failPendingPublishesOnRecovery's
// javadoc), it is just not reachable from this deterministic ordering, so it does not need a
// separate assertion here.
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
// sweep" racing a publish that registers while it is mid-flight.
@@ -109,8 +133,8 @@ class AmqpReplyInboxRecoveryRaceTest {
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
// for the registration rather than checking once: freshThread may still be contending for
// publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the
// sweep itself has already finished.
// publishChannelLock (behind the sweep's own hold on it, and possibly other stale threads
// still unwinding) even though the sweep itself has already finished.
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
int idx = -1;
while (idx < 0 && System.nanoTime() < deadline) {
@@ -42,7 +42,6 @@ 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;
@@ -580,83 +579,6 @@ 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