Skip to content

Commit a02732c

Browse files
authored
test(qwp): fix flaky port-bind race in WebSocket client tests (#47)
1 parent 38ab83c commit a02732c

21 files changed

Lines changed: 220 additions & 205 deletions

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/CleanShutdownNoReplayTest.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -73,12 +73,12 @@ public void testFullyAckedActiveDoesNotReplayAfterCleanRestart() throws Exceptio
7373
// Phase 1: server ACKs every frame. Sender writes a few rows,
7474
// flushes, then close() blocks for the default 5s drain — by the
7575
// time close returns, every frame has been ACK'd.
76-
int port1 = TestPorts.findUnusedPort();
7776
AckHandler ack1 = new AckHandler();
78-
try (TestWebSocketServer s1 = new TestWebSocketServer(port1, ack1)) {
77+
try (TestWebSocketServer s1 = new TestWebSocketServer(ack1)) {
7978
s1.start();
8079
Assert.assertTrue(s1.awaitStart(5, TimeUnit.SECONDS));
8180

81+
int port1 = s1.getPort();
8282
String cfg1 = "ws::addr=localhost:" + port1
8383
+ ";sf_dir=" + sfDir + ";";
8484
try (Sender sender = Sender.fromConfig(cfg1)) {
@@ -105,12 +105,12 @@ public void testFullyAckedActiveDoesNotReplayAfterCleanRestart() throws Exceptio
105105
// SAME slot dir. There is no unacked work — both rings should agree
106106
// there's nothing to send. The expected count of binary frames at
107107
// server 2 is zero.
108-
int port2 = port1 + 50;
109108
AckHandler ack2 = new AckHandler();
110-
try (TestWebSocketServer s2 = new TestWebSocketServer(port2, ack2)) {
109+
try (TestWebSocketServer s2 = new TestWebSocketServer(ack2)) {
111110
s2.start();
112111
Assert.assertTrue(s2.awaitStart(5, TimeUnit.SECONDS));
113112

113+
int port2 = s2.getPort();
114114
String cfg2 = "ws::addr=localhost:" + port2
115115
+ ";sf_dir=" + sfDir + ";";
116116
try (Sender ignored = Sender.fromConfig(cfg2)) {

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/CloseDrainDoubleSignalTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -103,12 +103,12 @@ public class CloseDrainDoubleSignalTest {
103103

104104
@Test(timeout = 30_000)
105105
public void testCloseDoesNotDoubleSignalWhenAsyncHandlerOwnsErrorAndDrainRuns() throws Exception {
106-
int port = TestPorts.findUnusedPort();
107106
GatedHaltHandler server = new GatedHaltHandler();
108-
try (TestWebSocketServer ws = new TestWebSocketServer(port, server)) {
107+
try (TestWebSocketServer ws = new TestWebSocketServer(server)) {
109108
ws.start();
110109
Assert.assertTrue(ws.awaitStart(5, TimeUnit.SECONDS));
111110

111+
int port = ws.getPort();
112112
// Memory mode + a positive drain timeout: drainOnClose() WILL run.
113113
String cfg = "ws::addr=localhost:" + port
114114
+ ";close_flush_timeout_millis=2000;";

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/CloseDrainTest.java

Lines changed: 20 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -54,13 +54,13 @@ public void testCloseBlocksUntilAckArrives() throws Exception {
5454
// Server delays every ACK by 800ms. With the default
5555
// close_flush_timeout_millis=60000, close() must wait for that
5656
// ACK before returning. Pre-fix close() returned within milliseconds.
57-
int port = TestPorts.findUnusedPort();
5857
long ackDelayMs = 800;
5958
DelayingAckHandler handler = new DelayingAckHandler(ackDelayMs);
60-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
59+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
6160
server.start();
6261
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
6362

63+
int port = server.getPort();
6464
String cfg = "ws::addr=localhost:" + port + ";"; // memory mode
6565
long elapsedMs;
6666
try (Sender sender = Sender.fromConfig(cfg)) {
@@ -82,13 +82,13 @@ public void testCloseFastWhenTimeoutIsZero() throws Exception {
8282
// Same delayed-ACK server, but with close_flush_timeout_millis=0
8383
// (fast close). close() must return immediately, well before the
8484
// ACK delay would have elapsed.
85-
int port = TestPorts.findUnusedPort();
8685
long ackDelayMs = 1500;
8786
DelayingAckHandler handler = new DelayingAckHandler(ackDelayMs);
88-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
87+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
8988
server.start();
9089
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
9190

91+
int port = server.getPort();
9292
String cfg = "ws::addr=localhost:" + port
9393
+ ";close_flush_timeout_millis=0;";
9494
long elapsedMs;
@@ -115,13 +115,13 @@ public void testCloseFastWhenTimeoutIsMinusOne() throws Exception {
115115
// sentinel in LineSenderBuilder, so the build path silently substitutes
116116
// DEFAULT_CLOSE_FLUSH_TIMEOUT_MILLIS (60s) and close() blocks for the
117117
// full ACK delay instead of returning fast.
118-
int port = TestPorts.findUnusedPort();
119118
long ackDelayMs = 1500;
120119
DelayingAckHandler handler = new DelayingAckHandler(ackDelayMs);
121-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
120+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
122121
server.start();
123122
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
124123

124+
int port = server.getPort();
125125
String cfg = "ws::addr=localhost:" + port
126126
+ ";close_flush_timeout_millis=-1;";
127127
long elapsedMs;
@@ -144,13 +144,13 @@ public void testCloseDrainTimesOutWhenAcksNeverArrive() throws Exception {
144144
// Server that buffers frames silently and never ACKs. close() must
145145
// throw a drain-timeout LineSenderException after roughly the
146146
// configured timeout — not hang forever and not return immediately.
147-
int port = TestPorts.findUnusedPort();
148147
long timeoutMs = 500;
149148
SilentHandler handler = new SilentHandler();
150-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
149+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
151150
server.start();
152151
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
153152

153+
int port = server.getPort();
154154
String cfg = "ws::addr=localhost:" + port
155155
+ ";close_flush_timeout_millis=" + timeoutMs + ";";
156156
long elapsedMs;
@@ -182,13 +182,13 @@ public void testDrainBlocksUntilAckArrivesAndReturnsTrue() throws Exception {
182182
// testCloseBlocksUntilAckArrives, but the wait happens inside the
183183
// explicit drain() call. The subsequent close() should be a near-
184184
// instant no-op because everything is already acked.
185-
int port = TestPorts.findUnusedPort();
186185
long ackDelayMs = 600;
187186
DelayingAckHandler handler = new DelayingAckHandler(ackDelayMs);
188-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
187+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
189188
server.start();
190189
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
191190

191+
int port = server.getPort();
192192
String cfg = "ws::addr=localhost:" + port + ";";
193193
try (Sender sender = Sender.fromConfig(cfg)) {
194194
sender.table("foo").longColumn("v", 1L).atNow();
@@ -225,13 +225,13 @@ public void testDrainBlocksUntilAckArrivesAndReturnsTrue() throws Exception {
225225
*/
226226
@Test
227227
public void testDrainAfterFlushWaitsForPriorUnackedFrames() throws Exception {
228-
int port = TestPorts.findUnusedPort();
229228
long ackDelayMs = 600;
230229
DelayingAckHandler handler = new DelayingAckHandler(ackDelayMs);
231-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
230+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
232231
server.start();
233232
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
234233

234+
int port = server.getPort();
235235
String cfg = "ws::addr=localhost:" + port + ";";
236236
try (Sender sender = Sender.fromConfig(cfg)) {
237237
sender.table("foo").longColumn("v", 1L).atNow();
@@ -259,12 +259,12 @@ public void testDrainReturnsFalseOnTimeoutAndSenderStillUsable() throws Exceptio
259259
// usable for further row writes after a false return; the
260260
// outstanding frames remain pending and close()'s own drain still
261261
// runs.
262-
int port = TestPorts.findUnusedPort();
263262
SilentHandler handler = new SilentHandler();
264-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
263+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
265264
server.start();
266265
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
267266

267+
int port = server.getPort();
268268
String cfg = "ws::addr=localhost:" + port + ";close_flush_timeout_millis=0;";
269269
try (Sender sender = Sender.fromConfig(cfg)) {
270270
sender.table("foo").longColumn("v", 1L).atNow();
@@ -287,12 +287,12 @@ public void testDrainNonZeroTimeoutOnFastServerReturnsImmediately() throws Excep
287287
// Fast server: every frame is acked promptly. drain(longTimeout)
288288
// must return true quickly -- no spurious wait when there is
289289
// nothing to wait for.
290-
int port = TestPorts.findUnusedPort();
291290
DelayingAckHandler handler = new DelayingAckHandler(0);
292-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
291+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
293292
server.start();
294293
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
295294

295+
int port = server.getPort();
296296
String cfg = "ws::addr=localhost:" + port + ";";
297297
try (Sender sender = Sender.fromConfig(cfg)) {
298298
sender.table("foo").longColumn("v", 1L).atNow();
@@ -307,9 +307,9 @@ public void testDrainNonZeroTimeoutOnFastServerReturnsImmediately() throws Excep
307307

308308
@Test
309309
public void testAsyncCloseDrainSucceedsWhenServerStartsDuringDrain() throws Exception {
310-
int port = TestPorts.findUnusedPort();
311310
DelayingAckHandler handler = new DelayingAckHandler(0);
312-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
311+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
312+
int port = server.getPort();
313313
String cfg = "ws::addr=localhost:" + port
314314
+ sfDirOpt()
315315
+ ";initial_connect_retry=async"
@@ -344,12 +344,12 @@ public void testAsyncCloseDrainSucceedsWhenServerStartsDuringDrain() throws Exce
344344

345345
@Test
346346
public void testAsyncCloseDrainSucceedsWhenServerWasUpAllAlong() throws Exception {
347-
int port = TestPorts.findUnusedPort();
348347
DelayingAckHandler handler = new DelayingAckHandler(0);
349-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
348+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
350349
server.start();
351350
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
352351

352+
int port = server.getPort();
353353
for (int i = 0; i < 20; i++) {
354354
String cfg = "ws::addr=localhost:" + port
355355
+ sfDirOpt()

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/CloseTerminalConflationTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -71,12 +71,12 @@ public class CloseTerminalConflationTest {
7171

7272
@Test(timeout = 30_000)
7373
public void testCloseSurfacesHaltAfterEarlierDropFlippedTheStickyFlag() throws Exception {
74-
int port = TestPorts.findUnusedPort();
7574
DropThenGatedHaltHandler server = new DropThenGatedHaltHandler();
76-
try (TestWebSocketServer ws = new TestWebSocketServer(port, server)) {
75+
try (TestWebSocketServer ws = new TestWebSocketServer(server)) {
7776
ws.start();
7877
Assert.assertTrue(ws.awaitStart(5, TimeUnit.SECONDS));
7978

79+
int port = ws.getPort();
8080
// Memory mode + a positive drain timeout: drainOnClose() WILL run.
8181
String cfg = "ws::addr=localhost:" + port
8282
+ ";close_flush_timeout_millis=2000;";

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/InitialConnectAsyncTest.java

Lines changed: 18 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -60,9 +60,9 @@ public void testAsyncAuthFailureDeliversToErrorInbox() throws Exception {
6060
// Server returns HTTP 401 on every upgrade attempt. Auth failures
6161
// are terminal at the I/O thread; in async mode they are
6262
// delivered as a SenderError, not thrown from fromConfig.
63-
int port = TestPorts.findUnusedPort();
64-
try (Always401Fixture fixture = new Always401Fixture(port)) {
63+
try (Always401Fixture fixture = new Always401Fixture()) {
6564
fixture.start();
65+
int port = fixture.getPort();
6666
ErrorInbox inbox = new ErrorInbox();
6767
String cfg = "ws::addr=localhost:" + port
6868
+ sfDirOpt() + ";initial_connect_retry=async"
@@ -159,9 +159,9 @@ public void testAsyncDeliversBufferedRowsWhenServerArrivesLate() {
159159
// appended to the cursor SF engine on the producer thread. The
160160
// I/O thread retries connect in the background; once the server
161161
// comes up, the buffered frame is sent and ACKed.
162-
int port = TestPorts.findUnusedPort();
163162
AckHandler handler = new AckHandler();
164-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
163+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
164+
int port = server.getPort();
165165
String cfg = "ws::addr=localhost:" + port
166166
+ sfDirOpt() + ";initial_connect_retry=async"
167167
+ ";reconnect_max_duration_millis=10000"
@@ -239,9 +239,9 @@ public void testConnectionLostBudgetExhaustionTagsDifferently() {
239239
// Because the loop did connect at least once before the outage,
240240
// the SenderError must use the connection-lost tag and the sender
241241
// must report wasEverConnected()==true.
242-
int port = TestPorts.findUnusedPort();
243242
AckHandler handler = new AckHandler();
244-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
243+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
244+
int port = server.getPort();
245245
server.start();
246246
Assert.assertTrue(server.awaitStart(5, java.util.concurrent.TimeUnit.SECONDS));
247247

@@ -299,8 +299,8 @@ public void testWasEverConnectedTrueImmediatelyInSyncMode() {
299299
// caller — there is no observable "never connected" window in
300300
// those modes, so misclassifying a budget exhaustion as
301301
// never-connected is impossible.
302-
int port = TestPorts.findUnusedPort();
303-
try (TestWebSocketServer server = new TestWebSocketServer(port, new AckHandler())) {
302+
try (TestWebSocketServer server = new TestWebSocketServer(new AckHandler())) {
303+
int port = server.getPort();
304304
server.start();
305305
Assert.assertTrue(server.awaitStart(5, java.util.concurrent.TimeUnit.SECONDS));
306306
String cfg = "ws::addr=localhost:" + port
@@ -446,8 +446,16 @@ private static class Always401Fixture implements AutoCloseable {
446446
private Thread acceptThread;
447447
private volatile boolean running;
448448

449-
Always401Fixture(int port) throws IOException {
450-
this.serverSocket = new ServerSocket(port);
449+
Always401Fixture() throws IOException {
450+
// Bind the listener up front on an OS-assigned loopback port and
451+
// hold it for the fixture's lifetime; read it back via getPort().
452+
// Owning the port from allocation to teardown avoids the bind race
453+
// a pre-selected port would carry.
454+
this.serverSocket = new ServerSocket(0, 50, java.net.InetAddress.getLoopbackAddress());
455+
}
456+
457+
int getPort() {
458+
return serverSocket.getLocalPort();
451459
}
452460

453461
@Override

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/InitialConnectRetryTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -100,14 +100,14 @@ public void testWithRetryGivesUpAfterCap() {
100100
}
101101

102102
@Test
103-
public void testWithRetrySucceedsWhenServerComesUpInTime() {
103+
public void testWithRetrySucceedsWhenServerComesUpInTime() throws Exception {
104104
// initial_connect_retry=true; we open the sender BEFORE starting
105105
// the server, then start the server in a background thread after
106106
// a short delay. The retry loop should see the server come up and
107107
// proceed cleanly.
108-
int port = TestPorts.findUnusedPort();
109108
AckHandler handler = new AckHandler();
110-
TestWebSocketServer server = new TestWebSocketServer(port, handler);
109+
TestWebSocketServer server = new TestWebSocketServer(handler);
110+
int port = server.getPort();
111111
Thread starter = new Thread(() -> {
112112
try {
113113
Thread.sleep(300);

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/IoThreadErrorSurfacedOnRowApiTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -64,9 +64,9 @@ public class IoThreadErrorSurfacedOnRowApiTest {
6464

6565
@Test
6666
public void testRowApiMethodSurfacesIoThreadTerminalError() throws Exception {
67-
int port = TestPorts.findUnusedPort();
6867
ErrorAckHandler handler = new ErrorAckHandler();
69-
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
68+
try (TestWebSocketServer server = new TestWebSocketServer(handler)) {
69+
int port = server.getPort();
7070
server.start();
7171
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
7272

core/src/test/java/io/questdb/client/test/cutlass/qwp/client/PrReviewRedTestsE2e.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -78,13 +78,13 @@ public class PrReviewRedTestsE2e {
7878
@Test
7979
public void testC4_handlerMustObserveTerminalErrorWhenInvoked() throws Exception {
8080
TestUtils.assertMemoryLeak(() -> {
81-
int port = TestPorts.findUnusedPort();
8281
int iterations = 30;
8382
AtomicInteger nullObservations = new AtomicInteger();
8483
AtomicInteger totalObservations = new AtomicInteger();
8584

8685
ParseErrorAckHandler serverHandler = new ParseErrorAckHandler();
87-
try (TestWebSocketServer server = new TestWebSocketServer(port, serverHandler)) {
86+
try (TestWebSocketServer server = new TestWebSocketServer(serverHandler)) {
87+
int port = server.getPort();
8888
server.start();
8989
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
9090

@@ -161,9 +161,9 @@ public void testC4_handlerMustObserveTerminalErrorWhenInvoked() throws Exception
161161
@Test
162162
public void testC11_postHaltFlushThrowsTypedLineSenderServerException() throws Exception {
163163
TestUtils.assertMemoryLeak(() -> {
164-
int port = TestPorts.findUnusedPort();
165164
ParseErrorAckHandler serverHandler = new ParseErrorAckHandler();
166-
try (TestWebSocketServer server = new TestWebSocketServer(port, serverHandler)) {
165+
try (TestWebSocketServer server = new TestWebSocketServer(serverHandler)) {
166+
int port = server.getPort();
167167
server.start();
168168
Assert.assertTrue(server.awaitStart(5, TimeUnit.SECONDS));
169169

0 commit comments

Comments
 (0)