Skip to content

Commit a5abbca

Browse files
authored
fix: preserve hard-strict priority when the worker pool is busy (#306)
1 parent 5326191 commit a5abbca

2 files changed

Lines changed: 35 additions & 2 deletions

File tree

‎rqueue-core/src/main/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPoller.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -208,7 +208,7 @@ void deactivate(int index, String queue, DeactivateType deactivateType) {
208208
if (deactivateType == DeactivateType.POLL_FAILED) {
209209
// Pause in case of connection errors or polling failures
210210
TimeoutUtils.sleepLog(backoffTime, false);
211-
} else {
211+
} else if (deactivateType == DeactivateType.NO_MESSAGE) {
212212
// Mark deactivation time if the queue is empty
213213
queueDeactivationTime.put(queue, System.currentTimeMillis());
214214
}

‎rqueue-core/src/test/java/com/github/sonus21/rqueue/listener/HardStrictPriorityPollerTest.java‎

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,23 @@
11
package com.github.sonus21.rqueue.listener;
22

3+
import static org.junit.jupiter.api.Assertions.assertFalse;
34
import static org.junit.jupiter.api.Assertions.assertTrue;
45
import static org.mockito.ArgumentMatchers.any;
56
import static org.mockito.ArgumentMatchers.anyInt;
67
import static org.mockito.ArgumentMatchers.anyString;
78
import static org.mockito.ArgumentMatchers.eq;
89
import static org.mockito.Mockito.lenient;
910
import static org.mockito.Mockito.mock;
11+
import static org.mockito.Mockito.never;
1012
import static org.mockito.Mockito.spy;
13+
import static org.mockito.Mockito.verify;
1114
import static org.mockito.Mockito.when;
1215

1316
import com.github.sonus21.TestBase;
1417
import com.github.sonus21.rqueue.CoreUnitTest;
1518
import com.github.sonus21.rqueue.core.RqueueBeanProvider;
1619
import com.github.sonus21.rqueue.listener.RqueueMessageListenerContainer.QueueStateMgr;
20+
import com.github.sonus21.rqueue.listener.RqueueMessagePoller.DeactivateType;
1721
import com.github.sonus21.rqueue.utils.Constants;
1822
import com.github.sonus21.rqueue.utils.QueueThreadPool;
1923
import com.github.sonus21.rqueue.utils.TimeoutUtils;
@@ -67,7 +71,7 @@ public void setUp() {
6771
rqueueBeanProvider,
6872
queueStateMgr,
6973
Collections.emptyList(),
70-
50L,
74+
60_000L, // Keep deactivation from expiring during assertions.
7175
50L,
7276
postProcessingHandler,
7377
new MessageHeaders(Collections.emptyMap()),
@@ -109,6 +113,35 @@ void testExistMessagesInHigherPriorityQueueReturnsTrue() {
109113
assertTrue(result, "Should return true because high priority queue has messages");
110114
}
111115

116+
@Test
117+
void testBusyPoolDoesNotHideHighPriorityMessages() throws Exception {
118+
QueueThreadPool busyPool = mock(QueueThreadPool.class);
119+
when(busyPool.acquire(1, poller.getSemaphoreWaitTime())).thenReturn(false);
120+
lenient().doReturn(true).when(poller).existAvailableMessagesForPoll(highDetail);
121+
lenient().doReturn(false).when(poller).existAvailableMessagesForPoll(lowDetail);
122+
123+
poller.poll(-1, highPriorityQueue, highDetail, busyPool);
124+
125+
verify(busyPool).acquire(1, poller.getSemaphoreWaitTime());
126+
assertTrue(
127+
poller.existMessagesInCurrentQueueOrHigherPriorityQueue(lowPriorityQueue, poller.queues),
128+
"A busy pool must not hide pending high-priority messages");
129+
verify(poller, never()).existAvailableMessagesForPoll(lowDetail);
130+
}
131+
132+
@Test
133+
void testEmptyQueueIsTemporarilySkipped() {
134+
lenient().doReturn(false).when(poller).existAvailableMessagesForPoll(lowDetail);
135+
136+
poller.deactivate(-1, highPriorityQueue, DeactivateType.NO_MESSAGE);
137+
138+
assertFalse(
139+
poller.existMessagesInCurrentQueueOrHigherPriorityQueue(lowPriorityQueue, poller.queues),
140+
"An empty high-priority queue should remain inactive during the polling interval");
141+
verify(poller, never()).existAvailableMessagesForPoll(highDetail);
142+
verify(poller).existAvailableMessagesForPoll(lowDetail);
143+
}
144+
112145
@Test
113146
void testStrictExecutionPreventsLowPriorityPoll() throws Exception {
114147
AtomicInteger highQueuePollCount = new AtomicInteger(0);

0 commit comments

Comments
 (0)