From 0e44571201fc4f2994f25c249573947f45b8546d Mon Sep 17 00:00:00 2001 From: 0x0fee <111068600+0x0fee@users.noreply.github.com> Date: Tue, 22 Sep 2026 20:58:53 +0300 Subject: [PATCH] fix: preserve hard-strict priority when the worker pool is busy --- .../listener/HardStrictPriorityPoller.java | 2 +- .../HardStrictPriorityPollerTest.java | 35 ++++++++++++++++++- 2 files changed, 35 insertions(+), 2 deletions(-) diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPoller.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPoller.java index 54386b2a..c29cfdea 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPoller.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPoller.java @@ -208,7 +208,7 @@ void deactivate(int index, String queue, DeactivateType deactivateType) { if (deactivateType == DeactivateType.POLL_FAILED) { // Pause in case of connection errors or polling failures TimeoutUtils.sleepLog(backoffTime, false); - } else { + } else if (deactivateType == DeactivateType.NO_MESSAGE) { // Mark deactivation time if the queue is empty queueDeactivationTime.put(queue, System.currentTimeMillis()); } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPollerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPollerTest.java index 5d9e0396..d59c1c50 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPollerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPollerTest.java @@ -1,5 +1,6 @@ package com.github.sonus21.rqueue.listener; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; @@ -7,13 +8,16 @@ import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import com.github.sonus21.TestBase; import com.github.sonus21.rqueue.CoreUnitTest; import com.github.sonus21.rqueue.core.RqueueBeanProvider; import com.github.sonus21.rqueue.listener.RqueueMessageListenerContainer.QueueStateMgr; +import com.github.sonus21.rqueue.listener.RqueueMessagePoller.DeactivateType; import com.github.sonus21.rqueue.utils.Constants; import com.github.sonus21.rqueue.utils.QueueThreadPool; import com.github.sonus21.rqueue.utils.TimeoutUtils; @@ -67,7 +71,7 @@ public void setUp() { rqueueBeanProvider, queueStateMgr, Collections.emptyList(), - 50L, + 60_000L, // Keep deactivation from expiring during assertions. 50L, postProcessingHandler, new MessageHeaders(Collections.emptyMap()), @@ -109,6 +113,35 @@ void testExistMessagesInHigherPriorityQueueReturnsTrue() { assertTrue(result, "Should return true because high priority queue has messages"); } + @Test + void testBusyPoolDoesNotHideHighPriorityMessages() throws Exception { + QueueThreadPool busyPool = mock(QueueThreadPool.class); + when(busyPool.acquire(1, poller.getSemaphoreWaitTime())).thenReturn(false); + lenient().doReturn(true).when(poller).existAvailableMessagesForPoll(highDetail); + lenient().doReturn(false).when(poller).existAvailableMessagesForPoll(lowDetail); + + poller.poll(-1, highPriorityQueue, highDetail, busyPool); + + verify(busyPool).acquire(1, poller.getSemaphoreWaitTime()); + assertTrue( + poller.existMessagesInCurrentQueueOrHigherPriorityQueue(lowPriorityQueue, poller.queues), + "A busy pool must not hide pending high-priority messages"); + verify(poller, never()).existAvailableMessagesForPoll(lowDetail); + } + + @Test + void testEmptyQueueIsTemporarilySkipped() { + lenient().doReturn(false).when(poller).existAvailableMessagesForPoll(lowDetail); + + poller.deactivate(-1, highPriorityQueue, DeactivateType.NO_MESSAGE); + + assertFalse( + poller.existMessagesInCurrentQueueOrHigherPriorityQueue(lowPriorityQueue, poller.queues), + "An empty high-priority queue should remain inactive during the polling interval"); + verify(poller, never()).existAvailableMessagesForPoll(highDetail); + verify(poller).existAvailableMessagesForPoll(lowDetail); + } + @Test void testStrictExecutionPreventsLowPriorityPoll() throws Exception { AtomicInteger highQueuePollCount = new AtomicInteger(0);