diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java index 4031ea6..97d0387 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/AmqpReplyInboxRecoveryRaceTest.java @@ -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) {