Skip to content

Commit 04b5385

Browse files
author
zileong.ong
committed
Fix pending queue lookup race
1 parent 2775db2 commit 04b5385

2 files changed

Lines changed: 39 additions & 4 deletions

File tree

‎src/main/java/io/nats/client/impl/NatsConsumer.java‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -83,15 +83,17 @@ public long getPendingByteLimit() {
8383
* {@link #setPendingLimits(long, long) setPendingLimits}.
8484
*/
8585
public long getPendingMessageCount() {
86-
return this.getMessageQueue() != null ? this.getMessageQueue().length() : 0;
86+
ConsumerMessageQueue queue = this.getMessageQueue();
87+
return queue == null ? 0 : queue.length();
8788
}
8889

8990
/**
9091
* @return the cumulative size of the messages waiting to be delivered/popped,
9192
* {@link #setPendingLimits(long, long) setPendingLimits}.
9293
*/
9394
public long getPendingByteCount() {
94-
return this.getMessageQueue() != null ? this.getMessageQueue().sizeInBytes() : 0;
95+
ConsumerMessageQueue queue = this.getMessageQueue();
96+
return queue == null ? 0 : queue.sizeInBytes();
9597
}
9698

9799
/**
@@ -249,4 +251,4 @@ public CompletableFuture<Boolean> drain(Duration timeout) throws InterruptedExce
249251
* Abstract method, called by the connection when the drain is complete.
250252
*/
251253
abstract void cleanUpAfterDrain();
252-
}
254+
}

‎src/test/java/io/nats/client/impl/SlowConsumerTests.java‎

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,11 +21,18 @@
2121
import java.util.concurrent.CompletableFuture;
2222
import java.util.concurrent.Future;
2323
import java.util.concurrent.TimeUnit;
24+
import java.util.concurrent.atomic.AtomicInteger;
2425

2526
import static org.junit.jupiter.api.Assertions.assertEquals;
2627

2728
public class SlowConsumerTests {
2829

30+
@Test
31+
public void testPendingCountsUseSingleQueueLookup() {
32+
assertEquals(0, new QueueInvalidatingConsumer().getPendingMessageCount());
33+
assertEquals(0, new QueueInvalidatingConsumer().getPendingByteCount());
34+
}
35+
2936
@Test
3037
public void testDefaultPendingLimits() throws Exception {
3138
try (NatsTestServer ts = new NatsTestServer(false);
@@ -234,4 +241,30 @@ public void testSlowSubscriberNotification() throws Exception {
234241
assertEquals(sub, slow.get(0));
235242
}
236243
}
237-
}
244+
245+
private static class QueueInvalidatingConsumer extends NatsConsumer {
246+
private final AtomicInteger queueLookups = new AtomicInteger();
247+
248+
QueueInvalidatingConsumer() {
249+
super(null);
250+
}
251+
252+
@Override
253+
public boolean isActive() {
254+
return true;
255+
}
256+
257+
@Override
258+
ConsumerMessageQueue getMessageQueue() {
259+
return queueLookups.getAndIncrement() == 0 ? new ConsumerMessageQueue() : null;
260+
}
261+
262+
@Override
263+
void sendUnsubForDrain() {
264+
}
265+
266+
@Override
267+
void cleanUpAfterDrain() {
268+
}
269+
}
270+
}

0 commit comments

Comments
 (0)