Skip to content

Commit 186c127

Browse files
committed
address review comments
1 parent 1e8db76 commit 186c127

6 files changed

Lines changed: 685 additions & 22 deletions

File tree

core/src/main/java/io/questdb/client/cutlass/qwp/client/QwpWebSocketSender.java

Lines changed: 7 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1757,10 +1757,10 @@ private void syncPing() {
17571757
if (received) {
17581758
if (sawBinaryAck) {
17591759
if (ackResponse.isDurableAck()) {
1760-
updateSyncDurableSeqTxns();
1760+
updateSyncSeqTxns(syncDurableSeqTxns);
17611761
} else if (ackResponse.isSuccess()) {
17621762
inFlightWindow.acknowledgeUpTo(ackResponse.getSequence());
1763-
updateSyncCommittedSeqTxns();
1763+
updateSyncSeqTxns(syncCommittedSeqTxns);
17641764
}
17651765
}
17661766
if (sawPong) {
@@ -1771,22 +1771,12 @@ private void syncPing() {
17711771
throw new LineSenderException("Ping timed out");
17721772
}
17731773

1774-
private void updateSyncCommittedSeqTxns() {
1774+
private void updateSyncSeqTxns(CharSequenceLongHashMap seqTxns) {
17751775
for (int i = 0, n = ackResponse.getTableEntryCount(); i < n; i++) {
17761776
String name = ackResponse.getTableName(i);
17771777
long seqTxn = ackResponse.getTableSeqTxn(i);
1778-
if (seqTxn > syncCommittedSeqTxns.get(name)) {
1779-
syncCommittedSeqTxns.put(name, seqTxn);
1780-
}
1781-
}
1782-
}
1783-
1784-
private void updateSyncDurableSeqTxns() {
1785-
for (int i = 0, n = ackResponse.getTableEntryCount(); i < n; i++) {
1786-
String name = ackResponse.getTableName(i);
1787-
long seqTxn = ackResponse.getTableSeqTxn(i);
1788-
if (seqTxn > syncDurableSeqTxns.get(name)) {
1789-
syncDurableSeqTxns.put(name, seqTxn);
1778+
if (seqTxn > seqTxns.get(name)) {
1779+
seqTxns.put(name, seqTxn);
17901780
}
17911781
}
17921782
}
@@ -1842,12 +1832,12 @@ private void waitForAck(long expectedSequence) {
18421832
if (ackResponse.isSuccess()) {
18431833
long sequence = ackResponse.getSequence();
18441834
inFlightWindow.acknowledgeUpTo(sequence);
1845-
updateSyncCommittedSeqTxns();
1835+
updateSyncSeqTxns(syncCommittedSeqTxns);
18461836
if (sequence >= expectedSequence) {
18471837
return;
18481838
}
18491839
} else if (ackResponse.isDurableAck()) {
1850-
updateSyncDurableSeqTxns();
1840+
updateSyncSeqTxns(syncDurableSeqTxns);
18511841
} else {
18521842
long sequence = ackResponse.getSequence();
18531843
String errorMessage = ackResponse.getErrorMessage();

core/src/main/java/io/questdb/client/cutlass/qwp/client/WebSocketSendQueue.java

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -809,11 +809,8 @@ private static void advanceSeqTxn(ConcurrentHashMap<String, SeqTxn> map, String
809809
SeqTxn existing = map.get(tableName);
810810
if (existing != null) {
811811
existing.advance(seqTxn);
812-
} else {
813-
existing = map.putIfAbsent(tableName, new SeqTxn(seqTxn));
814-
if (existing != null) {
815-
existing.advance(seqTxn);
816-
}
812+
return;
817813
}
814+
map.computeIfAbsent(tableName, k -> new SeqTxn(seqTxn)).advance(seqTxn);
818815
}
819816
}

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

Lines changed: 241 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424

2525
package io.questdb.client.test.cutlass.qwp.client;
2626

27+
import io.questdb.client.cutlass.line.LineSenderException;
2728
import io.questdb.client.cutlass.qwp.client.QwpWebSocketSender;
2829
import io.questdb.client.cutlass.qwp.client.WebSocketResponse;
2930
import io.questdb.client.cutlass.qwp.websocket.WebSocketCloseCode;
@@ -34,8 +35,13 @@
3435
import org.junit.Test;
3536

3637
import java.io.IOException;
38+
import java.io.InputStream;
39+
import java.net.ServerSocket;
40+
import java.net.Socket;
41+
import java.nio.charset.StandardCharsets;
3742
import java.util.concurrent.TimeUnit;
3843
import java.util.concurrent.atomic.AtomicLong;
44+
import java.util.concurrent.atomic.AtomicReference;
3945

4046
/**
4147
* Integration tests for QWP v1 WebSocket ACK delivery mechanism.
@@ -195,6 +201,151 @@ public void testSyncFlushIgnoresPingAndWaitsForAck() throws Exception {
195201
}
196202
}
197203

204+
@Test
205+
public void testDurableAckUpgradeHeaderNotSentByDefault() throws Exception {
206+
int port = TEST_PORT + 31;
207+
AtomicReference<String> capturedRequest = new AtomicReference<>();
208+
209+
try (ServerSocket serverSocket = new ServerSocket(port)) {
210+
serverSocket.setSoTimeout(5000);
211+
212+
Thread serverThread = new Thread(() -> {
213+
try {
214+
Socket client = serverSocket.accept();
215+
InputStream in = client.getInputStream();
216+
StringBuilder request = new StringBuilder();
217+
byte[] buf = new byte[1];
218+
while (true) {
219+
int read = in.read(buf);
220+
if (read <= 0) {
221+
break;
222+
}
223+
request.append((char) buf[0]);
224+
if (request.toString().endsWith("\r\n\r\n")) {
225+
break;
226+
}
227+
}
228+
capturedRequest.set(request.toString());
229+
client.close();
230+
} catch (Exception e) {
231+
// expected
232+
}
233+
});
234+
serverThread.start();
235+
236+
try {
237+
QwpWebSocketSender.connect("localhost", port, null,
238+
0, 0, 0, 1, null).close();
239+
} catch (LineSenderException e) {
240+
// expected - server doesn't complete handshake
241+
}
242+
243+
serverThread.join(5000);
244+
245+
String request = capturedRequest.get();
246+
Assert.assertNotNull("Server should have received upgrade request", request);
247+
Assert.assertFalse("Request should NOT contain X-QWP-Request-Durable-Ack header",
248+
request.contains("X-QWP-Request-Durable-Ack"));
249+
}
250+
}
251+
252+
@Test
253+
public void testDurableAckUpgradeHeaderSent() throws Exception {
254+
int port = TEST_PORT + 30;
255+
AtomicReference<String> capturedRequest = new AtomicReference<>();
256+
257+
try (ServerSocket serverSocket = new ServerSocket(port)) {
258+
serverSocket.setSoTimeout(5000);
259+
260+
Thread serverThread = new Thread(() -> {
261+
try {
262+
Socket client = serverSocket.accept();
263+
InputStream in = client.getInputStream();
264+
StringBuilder request = new StringBuilder();
265+
byte[] buf = new byte[1];
266+
while (true) {
267+
int read = in.read(buf);
268+
if (read <= 0) {
269+
break;
270+
}
271+
request.append((char) buf[0]);
272+
if (request.toString().endsWith("\r\n\r\n")) {
273+
break;
274+
}
275+
}
276+
capturedRequest.set(request.toString());
277+
client.close();
278+
} catch (Exception e) {
279+
// expected
280+
}
281+
});
282+
serverThread.start();
283+
284+
try {
285+
QwpWebSocketSender.connect("localhost", port, null,
286+
0, 0, 0, 1, null,
287+
QwpWebSocketSender.DEFAULT_MAX_SCHEMAS_PER_CONNECTION,
288+
true).close();
289+
} catch (LineSenderException e) {
290+
// expected - server doesn't complete handshake
291+
}
292+
293+
serverThread.join(5000);
294+
295+
String request = capturedRequest.get();
296+
Assert.assertNotNull("Server should have received upgrade request", request);
297+
Assert.assertTrue("Request should contain X-QWP-Request-Durable-Ack header",
298+
request.contains("X-QWP-Request-Durable-Ack: true"));
299+
}
300+
}
301+
302+
@Test
303+
public void testSyncDurableAckDuringWaitForAck() throws Exception {
304+
int port = TEST_PORT + 25;
305+
DurableAckThenStatusOkHandler handler = new DurableAckThenStatusOkHandler();
306+
307+
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
308+
server.start();
309+
Assert.assertTrue("Server failed to start", server.awaitStart(5, TimeUnit.SECONDS));
310+
311+
// window=1 for sync mode
312+
try (QwpWebSocketSender sender = QwpWebSocketSender.connect(
313+
"localhost", port, null, 0, 0, 0, 1, null)) {
314+
sender.table("trades")
315+
.longColumn("price", 100)
316+
.atNow();
317+
sender.flush();
318+
319+
Assert.assertEquals(42L, sender.getHighestDurableSeqTxn("trades"));
320+
Assert.assertEquals(10L, sender.getHighestAckedSeqTxn("trades"));
321+
}
322+
}
323+
}
324+
325+
@Test
326+
public void testSyncFlushUpdatesCommittedSeqTxnsWithTableEntries() throws Exception {
327+
int port = TEST_PORT + 24;
328+
AckWithTableEntriesHandler handler = new AckWithTableEntriesHandler();
329+
330+
try (TestWebSocketServer server = new TestWebSocketServer(port, handler)) {
331+
server.start();
332+
Assert.assertTrue("Server failed to start", server.awaitStart(5, TimeUnit.SECONDS));
333+
334+
// window=1 for sync mode
335+
try (QwpWebSocketSender sender = QwpWebSocketSender.connect(
336+
"localhost", port, null, 0, 0, 0, 1, null)) {
337+
sender.table("trades")
338+
.longColumn("price", 100)
339+
.atNow();
340+
sender.flush();
341+
342+
Assert.assertEquals(10L, sender.getHighestAckedSeqTxn("trades"));
343+
Assert.assertEquals(20L, sender.getHighestAckedSeqTxn("orders"));
344+
Assert.assertEquals(-1L, sender.getHighestAckedSeqTxn("other"));
345+
}
346+
}
347+
}
348+
198349
/**
199350
* Creates a binary ACK response using WebSocketResponse format.
200351
* Format: status (1) + sequence (8) + tableCount (2, zero entries)
@@ -220,6 +371,75 @@ private static byte[] createAckResponse(long sequence) {
220371
return response;
221372
}
222373

374+
private static byte[] createAckResponseWithTables(long sequence, String[] tableNames, long[] seqTxns) {
375+
byte[][] nameBytes = new byte[tableNames.length][];
376+
int size = 1 + 8 + 2;
377+
for (int i = 0; i < tableNames.length; i++) {
378+
nameBytes[i] = tableNames[i].getBytes(StandardCharsets.UTF_8);
379+
size += 2 + nameBytes[i].length + 8;
380+
}
381+
382+
byte[] response = new byte[size];
383+
int offset = 0;
384+
response[offset++] = WebSocketResponse.STATUS_OK;
385+
for (int i = 0; i < 8; i++) {
386+
response[offset++] = (byte) ((sequence >> (i * 8)) & 0xFF);
387+
}
388+
response[offset++] = (byte) (tableNames.length & 0xFF);
389+
response[offset++] = (byte) ((tableNames.length >> 8) & 0xFF);
390+
for (int i = 0; i < tableNames.length; i++) {
391+
response[offset++] = (byte) (nameBytes[i].length & 0xFF);
392+
response[offset++] = (byte) ((nameBytes[i].length >> 8) & 0xFF);
393+
System.arraycopy(nameBytes[i], 0, response, offset, nameBytes[i].length);
394+
offset += nameBytes[i].length;
395+
for (int j = 0; j < 8; j++) {
396+
response[offset++] = (byte) ((seqTxns[i] >> (j * 8)) & 0xFF);
397+
}
398+
}
399+
return response;
400+
}
401+
402+
private static byte[] createDurableAckResponse(String[] tableNames, long[] seqTxns) {
403+
byte[][] nameBytes = new byte[tableNames.length][];
404+
int size = 1 + 2;
405+
for (int i = 0; i < tableNames.length; i++) {
406+
nameBytes[i] = tableNames[i].getBytes(StandardCharsets.UTF_8);
407+
size += 2 + nameBytes[i].length + 8;
408+
}
409+
410+
byte[] response = new byte[size];
411+
int offset = 0;
412+
response[offset++] = WebSocketResponse.STATUS_DURABLE_ACK;
413+
response[offset++] = (byte) (tableNames.length & 0xFF);
414+
response[offset++] = (byte) ((tableNames.length >> 8) & 0xFF);
415+
for (int i = 0; i < tableNames.length; i++) {
416+
response[offset++] = (byte) (nameBytes[i].length & 0xFF);
417+
response[offset++] = (byte) ((nameBytes[i].length >> 8) & 0xFF);
418+
System.arraycopy(nameBytes[i], 0, response, offset, nameBytes[i].length);
419+
offset += nameBytes[i].length;
420+
for (int j = 0; j < 8; j++) {
421+
response[offset++] = (byte) ((seqTxns[i] >> (j * 8)) & 0xFF);
422+
}
423+
}
424+
return response;
425+
}
426+
427+
private static class AckWithTableEntriesHandler implements TestWebSocketServer.WebSocketServerHandler {
428+
private final AtomicLong nextSequence = new AtomicLong(0);
429+
430+
@Override
431+
public void onBinaryMessage(TestWebSocketServer.ClientHandler client, byte[] data) {
432+
long sequence = nextSequence.getAndIncrement();
433+
try {
434+
client.sendBinary(createAckResponseWithTables(sequence,
435+
new String[]{"trades", "orders"},
436+
new long[]{10L, 20L}));
437+
} catch (IOException e) {
438+
LOG.error("Failed to send ACK with tables", e);
439+
}
440+
}
441+
}
442+
223443
private static class ClosingServerHandler implements TestWebSocketServer.WebSocketServerHandler {
224444
@Override
225445
public void onBinaryMessage(TestWebSocketServer.ClientHandler client, byte[] data) {
@@ -261,6 +481,27 @@ public void onBinaryMessage(TestWebSocketServer.ClientHandler client, byte[] dat
261481
}
262482
}
263483

484+
private static class DurableAckThenStatusOkHandler implements TestWebSocketServer.WebSocketServerHandler {
485+
private final AtomicLong nextSequence = new AtomicLong(0);
486+
487+
@Override
488+
public void onBinaryMessage(TestWebSocketServer.ClientHandler client, byte[] data) {
489+
long sequence = nextSequence.getAndIncrement();
490+
try {
491+
// Send durable ACK first
492+
client.sendBinary(createDurableAckResponse(
493+
new String[]{"trades"},
494+
new long[]{42L}));
495+
// Then send STATUS_OK with committed seqTxns
496+
client.sendBinary(createAckResponseWithTables(sequence,
497+
new String[]{"trades"},
498+
new long[]{10L}));
499+
} catch (IOException e) {
500+
LOG.error("Failed to send ACK frames", e);
501+
}
502+
}
503+
}
504+
264505
private static class InvalidAckPayloadHandler implements TestWebSocketServer.WebSocketServerHandler {
265506
@Override
266507
public void onBinaryMessage(TestWebSocketServer.ClientHandler client, byte[] data) {

0 commit comments

Comments
 (0)