Skip to content

Commit dc2897d

Browse files
authored
Merge branch 'main' into fix-regression-getStreamTimeout
2 parents f8cc051 + 830c46a commit dc2897d

3 files changed

Lines changed: 41 additions & 53 deletions

File tree

‎src/test/java/io/nats/client/ConnectTests.java‎

Lines changed: 0 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,6 @@
1818
import io.nats.client.api.ServerInfo;
1919
import io.nats.client.impl.ListenerForTesting;
2020
import io.nats.client.impl.SimulateSocketDataPortException;
21-
import io.nats.client.utils.TestBase;
2221
import org.junit.jupiter.api.Test;
2322

2423
import java.io.IOException;
@@ -30,7 +29,6 @@
3029
import java.util.concurrent.CountDownLatch;
3130
import java.util.concurrent.TimeUnit;
3231
import java.util.concurrent.atomic.AtomicBoolean;
33-
import java.util.concurrent.atomic.AtomicLong;
3432

3533
import static io.nats.client.utils.TestBase.*;
3634
import static org.junit.jupiter.api.Assertions.*;
@@ -643,42 +641,4 @@ void testConnectWithHappyEyeballsShortCircuitCoverage() throws Exception {
643641
}
644642
}
645643

646-
@Test
647-
void testConnectPendingCountCoverage() throws Exception {
648-
TestBase.runInJsServer(nc -> {
649-
AtomicLong outgoingPendingMessageCount = new AtomicLong();
650-
AtomicLong outgoingPendingBytes = new AtomicLong();
651-
652-
AtomicBoolean tKeepGoing = new AtomicBoolean(true);
653-
Thread t = new Thread(() -> {
654-
while (tKeepGoing.get()) {
655-
outgoingPendingMessageCount.set(Math.max(outgoingPendingMessageCount.get(), nc.outgoingPendingMessageCount()));
656-
outgoingPendingBytes.set(Math.max(outgoingPendingBytes.get(), nc.outgoingPendingBytes()));
657-
try {
658-
Thread.sleep(10);
659-
}
660-
catch (InterruptedException e) {
661-
throw new RuntimeException(e);
662-
}
663-
}
664-
});
665-
t.start();
666-
667-
String subject = subject();
668-
// Keep the total queued bytes below the reconnect buffer size (default 8MB).
669-
// This test only needs the outgoing-pending counters to report a backlog;
670-
// flooding past the reconnect buffer would (correctly) throw if the connection
671-
// briefly reconnects under load on a slow/contended machine, which is not what
672-
// this test targets.
673-
byte[] data = new byte[2 * 1024];
674-
for (int x = 0; x < 3000; x++) {
675-
nc.publish(subject, data);
676-
}
677-
tKeepGoing.set(false);
678-
t.join();
679-
680-
assertTrue(outgoingPendingMessageCount.get() > 0);
681-
assertTrue(outgoingPendingBytes.get() > outgoingPendingMessageCount.get() * 1000);
682-
});
683-
}
684644
}

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

Lines changed: 17 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2035,14 +2035,17 @@ else if (mcount == 4) {
20352035
StreamInfo si = jsm.getStreamInfo(rawStream);
20362036
assertEquals(1, si.getStreamState().getMsgCount());
20372037

2038-
long safety = 0;
2038+
// Bound the wait by time, not by a poll count. How much wall time a fixed number
2039+
// of unthrottled getStreamInfo round trips buys depends entirely on the machine.
2040+
long timeoutAt = System.currentTimeMillis() + 10_000;
20392041
long gotZero = -1;
2040-
while (++safety < 10000 && errorLatch.getCount() > 0) {
2042+
while (System.currentTimeMillis() < timeoutAt && errorLatch.getCount() > 0) {
20412043
si = jsm.getStreamInfo(rawStream);
20422044
if (si.getStreamState().getMsgCount() == 0) {
20432045
gotZero = System.currentTimeMillis();
20442046
break;
20452047
}
2048+
sleep(10);
20462049
}
20472050
assertEquals(1, errorLatch.getCount(), error.get());
20482051
assertEquals(2, messages.get());
@@ -2058,14 +2061,15 @@ else if (mcount == 4) {
20582061
si = jsm.getStreamInfo(rawStream);
20592062
assertEquals(1, si.getStreamState().getMsgCount());
20602063

2061-
safety = 0;
2064+
timeoutAt = System.currentTimeMillis() + 10_000;
20622065
gotZero = -1;
2063-
while (++safety < 10000 && errorLatch.getCount() > 0) {
2066+
while (System.currentTimeMillis() < timeoutAt && errorLatch.getCount() > 0) {
20642067
si = jsm.getStreamInfo(rawStream);
20652068
if (si.getStreamState().getMsgCount() == 0) {
20662069
gotZero = System.currentTimeMillis();
20672070
break;
20682071
}
2072+
sleep(10);
20692073
}
20702074
assertEquals(1, errorLatch.getCount(), error.get());
20712075
assertEquals(4, messages.get());
@@ -2141,19 +2145,22 @@ else if (mcount == 4) {
21412145
StreamInfo si = jsm.getStreamInfo(rawStream);
21422146
assertEquals(1, si.getStreamState().getMsgCount());
21432147

2144-
kv.delete(key);
21452148
long mark = System.currentTimeMillis();
2149+
kv.delete(key);
21462150
si = jsm.getStreamInfo(rawStream);
21472151
assertEquals(1, si.getStreamState().getMsgCount());
21482152

2149-
long safety = 0;
2153+
// Bound the wait by time, not by a poll count. How much wall time a fixed number
2154+
// of unthrottled getStreamInfo round trips buys depends entirely on the machine.
2155+
long timeoutAt = System.currentTimeMillis() + 10_000;
21502156
long gotZero = -1;
2151-
while (++safety < 10000 && errorLatch.getCount() > 0) {
2157+
while (System.currentTimeMillis() < timeoutAt && errorLatch.getCount() > 0) {
21522158
si = jsm.getStreamInfo(rawStream);
21532159
if (si.getStreamState().getMsgCount() == 0) {
21542160
gotZero = System.currentTimeMillis();
21552161
break;
21562162
}
2163+
sleep(10);
21572164
}
21582165
assertEquals(1, errorLatch.getCount(), error.get());
21592166
assertEquals(2, messages.get());
@@ -2165,14 +2172,15 @@ else if (mcount == 4) {
21652172
mark = System.currentTimeMillis();
21662173
kv.purge(key);
21672174

2168-
safety = 0;
2175+
timeoutAt = System.currentTimeMillis() + 10_000;
21692176
gotZero = -1;
2170-
while (++safety < 10000 && errorLatch.getCount() > 0) {
2177+
while (System.currentTimeMillis() < timeoutAt && errorLatch.getCount() > 0) {
21712178
si = jsm.getStreamInfo(rawStream);
21722179
if (si.getStreamState().getMsgCount() == 0) {
21732180
gotZero = System.currentTimeMillis();
21742181
break;
21752182
}
2183+
sleep(10);
21762184
}
21772185
assertEquals(1, errorLatch.getCount(), error.get());
21782186
assertEquals(4, messages.get());

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

Lines changed: 24 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -18,10 +18,7 @@
1818
import org.junit.jupiter.api.Test;
1919

2020
import java.io.IOException;
21-
import java.util.concurrent.ExecutorService;
22-
import java.util.concurrent.Executors;
23-
import java.util.concurrent.ScheduledExecutorService;
24-
import java.util.concurrent.ThreadFactory;
21+
import java.util.concurrent.*;
2522
import java.util.concurrent.atomic.AtomicLong;
2623

2724
import static io.nats.client.utils.TestBase.*;
@@ -323,4 +320,27 @@ public void testReaderWriterThreadFactories() throws Exception {
323320
assertFalse(writerFactory.created.isEmpty());
324321
}
325322
}
323+
324+
@Test
325+
public void testOutgoingPendingCountCoverage() throws Exception {
326+
runInServer(nc -> {
327+
NatsConnection conn = (NatsConnection)nc;
328+
329+
// Stop the writer so nothing drains the outgoing queue, giving a deterministic backlog
330+
// to exercise the pending-count getters against. Reading them while the writer is live
331+
// is an unwinnable race (a fast machine drains to 0; a slow machine backs up past the
332+
// reconnect buffer and the publish throws), and they can't be read during reconnect at
333+
// all because that path holds closeSocketLock for the whole reconnect.
334+
conn.getWriter().stop().get(LONG_TIMEOUT_MS, TimeUnit.MILLISECONDS);
335+
336+
String subject = subject();
337+
byte[] data = new byte[2 * 1024]; // > 1000 bytes so pending bytes > count * 1000
338+
for (int x = 0; x < 20; x++) {
339+
conn.publish(subject, data);
340+
}
341+
342+
assertTrue(conn.outgoingPendingMessageCount() > 0);
343+
assertTrue(conn.outgoingPendingBytes() > conn.outgoingPendingMessageCount() * 1000);
344+
});
345+
}
326346
}

0 commit comments

Comments
 (0)