diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/WebSocketResponse.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/WebSocketResponse.java index 81d59d28..404f5eb2 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/WebSocketResponse.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/WebSocketResponse.java @@ -101,6 +101,7 @@ public class WebSocketResponse { // subsequent sight is allocation-free. private final Utf8SequenceObjHashMap tableNameCache = new Utf8SequenceObjHashMap<>(); private final ObjList tableNames = new ObjList<>(); + private final ObjList tableDirNames = new ObjList<>(); private final LongList tableSeqTxns = new LongList(); private String errorMessage; private int errorMessageUtf8Length; @@ -118,11 +119,12 @@ public WebSocketResponse() { * Creates a durable-upload ACK response with a single table entry. */ @TestOnly - public static WebSocketResponse durableAck(String tableName, long seqTxn) { + public static WebSocketResponse durableAck(String tableName, String tableDirName, long seqTxn) { WebSocketResponse response = new WebSocketResponse(); response.status = STATUS_DURABLE_ACK; response.sequence = -1; response.tableNames.add(tableName); + response.tableDirNames.add(tableDirName); response.tableSeqTxns.add(seqTxn); return response; } @@ -164,7 +166,7 @@ public static boolean isStructurallyValid(long ptr, int length) { if (length < MIN_DURABLE_ACK_SIZE) { return false; } - return validateTableEntries(ptr + 1, length - 1); + return validateDurableAckEntries(ptr + 1, length - 1); } // Error response @@ -245,6 +247,20 @@ public String getTableName(int index) { return tableNames.getQuick(index); } + /** + * Returns the server-side table directory name that accompanied the + * durable-ack entry at {@code index}. Used by the sender loop as an + * incarnation discriminator: when the dir name for a table name changes, + * the table was dropped and re-created on the server, so per-table + * durable-upload watermarks must be reset rather than max-merged. + * + * @param index entry index + * @return the table dir name, or {@code null} for non-durable-ack frames + */ + public String getTableDirName(int index) { + return tableDirNames.size() > index ? tableDirNames.getQuick(index) : null; + } + public long getTableSeqTxn(int index) { return tableSeqTxns.getQuick(index); } @@ -272,6 +288,7 @@ public boolean isSuccess() { */ public boolean readFrom(long ptr, int length) { tableNames.clear(); + tableDirNames.clear(); tableSeqTxns.clear(); if (length < 1) { @@ -297,7 +314,9 @@ public boolean readFrom(long ptr, int length) { sequence = -1; errorMessage = null; errorMessageUtf8Length = -1; - return readTableEntries(ptr + 1, length - 1); + // Durable ack entries carry the table dir name as an incarnation + // discriminator: [nameLen(2) + name(N) + dirLen(2) + dir(M) + seqTxn(8)] + return readDurableAckEntries(ptr + 1, length - 1); } // Error response @@ -334,7 +353,7 @@ public int serializedSize() { return MIN_OK_RESPONSE_SIZE + tableEntriesSize(); } if (status == STATUS_DURABLE_ACK) { - return MIN_DURABLE_ACK_SIZE + tableEntriesSize(); + return MIN_DURABLE_ACK_SIZE + durableAckEntriesSize(); } return MIN_ERROR_RESPONSE_SIZE + getErrorMessageUtf8Length(); } @@ -369,7 +388,7 @@ public int writeTo(long ptr) { offset += 8; offset += writeTableEntries(ptr + offset); } else if (status == STATUS_DURABLE_ACK) { - offset += writeTableEntries(ptr + offset); + offset += writeDurableAckEntries(ptr + offset); } else { Unsafe.getUnsafe().putLong(ptr + offset, sequence); offset += 8; @@ -420,6 +439,53 @@ private boolean readTableEntries(long ptr, int remaining) { return remaining == offset; } + /** + * Reads durable-ack table entries that carry the table dir name as an + * incarnation discriminator. + *

+ * Format: [nameLen(2) + nameUtf8(N) + dirLen(2) + dirUtf8(M) + seqTxn(8)] * count + */ + private boolean readDurableAckEntries(long ptr, int remaining) { + if (remaining < 2) { + return false; + } + int tableCount = Unsafe.getUnsafe().getShort(ptr) & 0xFFFF; + int offset = 2; + for (int i = 0; i < tableCount; i++) { + // Table name + if (remaining < offset + 2) { + return false; + } + int nameLen = Unsafe.getUnsafe().getShort(ptr + offset) & 0xFFFF; + offset += 2; + if (nameLen == 0 || remaining < offset + nameLen + 2) { + return false; + } + long nameLo = ptr + offset; + long nameHi = nameLo + nameLen; + offset += nameLen; + // Dir name (incarnation discriminator) + if (remaining < offset + 2) { + return false; + } + int dirLen = Unsafe.getUnsafe().getShort(ptr + offset) & 0xFFFF; + offset += 2; + if (dirLen == 0 || remaining < offset + dirLen + 8) { + return false; + } + long dirLo = ptr + offset; + long dirHi = dirLo + dirLen; + offset += dirLen; + // SeqTxn + long seqTxn = Unsafe.getUnsafe().getLong(ptr + offset); + offset += 8; + tableNames.add(internTableName(nameLo, nameHi)); + tableDirNames.add(Utf8s.stringFromUtf8Bytes(dirLo, dirHi)); + tableSeqTxns.add(seqTxn); + } + return remaining == offset; + } + private String internTableName(long lo, long hi) { lookupKey.of(lo, hi); int keyIndex = tableNameCache.keyIndex(lookupKey); @@ -462,6 +528,36 @@ private static boolean validateTableEntries(long ptr, int remaining) { return remaining == offset; } + private static boolean validateDurableAckEntries(long ptr, int remaining) { + if (remaining < 2) { + return false; + } + int tableCount = Unsafe.getUnsafe().getShort(ptr) & 0xFFFF; + int offset = 2; + for (int i = 0; i < tableCount; i++) { + if (remaining < offset + 2) { + return false; + } + int nameLen = Unsafe.getUnsafe().getShort(ptr + offset) & 0xFFFF; + offset += 2; + if (nameLen == 0 || remaining < offset + nameLen + 2) { + return false; + } + offset += nameLen; + // Dir name + if (remaining < offset + 2) { + return false; + } + int dirLen = Unsafe.getUnsafe().getShort(ptr + offset) & 0xFFFF; + offset += 2; + if (dirLen == 0 || remaining < offset + dirLen + 8) { + return false; + } + offset += dirLen + 8; + } + return remaining == offset; + } + private int writeTableEntries(long ptr) { int offset = 0; int count = tableNames.size(); @@ -481,6 +577,45 @@ private int writeTableEntries(long ptr) { return offset; } + private int writeDurableAckEntries(long ptr) { + int offset = 0; + int count = tableNames.size(); + Unsafe.getUnsafe().putShort(ptr + offset, (short) count); + offset += 2; + for (int i = 0; i < count; i++) { + // Table name + byte[] nameBytes = tableNames.getQuick(i).getBytes(StandardCharsets.UTF_8); + Unsafe.getUnsafe().putShort(ptr + offset, (short) nameBytes.length); + offset += 2; + for (int j = 0; j < nameBytes.length; j++) { + Unsafe.getUnsafe().putByte(ptr + offset + j, nameBytes[j]); + } + offset += nameBytes.length; + // Dir name (incarnation discriminator) + byte[] dirBytes = tableDirNames.getQuick(i).getBytes(StandardCharsets.UTF_8); + Unsafe.getUnsafe().putShort(ptr + offset, (short) dirBytes.length); + offset += 2; + for (int j = 0; j < dirBytes.length; j++) { + Unsafe.getUnsafe().putByte(ptr + offset + j, dirBytes[j]); + } + offset += dirBytes.length; + // SeqTxn + Unsafe.getUnsafe().putLong(ptr + offset, tableSeqTxns.getQuick(i)); + offset += 8; + } + return offset; + } + + private int durableAckEntriesSize() { + int size = 0; + for (int i = 0, n = tableNames.size(); i < n; i++) { + size += 2 + tableNames.getQuick(i).getBytes(StandardCharsets.UTF_8).length + + 2 + tableDirNames.getQuick(i).getBytes(StandardCharsets.UTF_8).length + + 8; + } + return size; + } + private int getErrorMessageUtf8Length() { if (status == STATUS_OK || status == STATUS_DURABLE_ACK || errorMessage == null || errorMessage.isEmpty()) { errorMessageUtf8Length = 0; diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java index 5a66c302..f97d9fb0 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/CursorWebSocketSendLoop.java @@ -262,6 +262,12 @@ public final class CursorWebSocketSendLoop implements QuietCloseable { // by the server -- holding stale watermarks across the wire boundary // would falsely advance trim before re-confirmation. private final CharSequenceLongHashMap durableTableWatermarks = new CharSequenceLongHashMap(); + // Per-table dir name tracking for incarnation change detection. + // Updated from STATUS_DURABLE_ACK frame entries alongside the watermark. + // When the dir name for a table changes (drop/recreate), the watermark + // is reset for that table since the new incarnation starts from a low + // seqTxn that must not be covered by the old incarnation's high watermark. + private final java.util.HashMap durableTableDirNames = new java.util.HashMap<>(); // Pre-converted to nanos. Consulted only by the orphan terminal policy. Zero disables // the dwell entirely (count-only escalation at MAX_CATCHUP_CAP_GAP_ATTEMPTS); the // user-facing 5-minute default is applied at the config layer. @@ -1630,7 +1636,18 @@ private void applyDurableAck() { int n = response.getTableEntryCount(); for (int i = 0; i < n; i++) { String name = response.getTableName(i); + String dirName = response.getTableDirName(i); long seqTxn = response.getTableSeqTxn(i); + // Incarnation change detection: when the table's dir name differs + // from what we last saw on this connection, the table was dropped + // and re-created under the same name. The old incarnation's high + // watermark must not cover the new incarnation's low seqTxns, so + // reset this table's watermark before applying the new value. + String lastDirName = durableTableDirNames.get(name); + if (lastDirName != null && !lastDirName.equals(dirName)) { + durableTableWatermarks.put(name, -1L); + } + durableTableDirNames.put(name, dirName); long current = durableTableWatermarks.get(name); if (seqTxn > current) { durableTableWatermarks.put(name, seqTxn); @@ -1663,6 +1680,7 @@ private void clearDurableAckTracking() { releasePendingEntry(pendingDurable.pollFirst()); } durableTableWatermarks.clear(); + durableTableDirNames.clear(); // Reset the keepalive throttle so the new connection can prod the // server immediately rather than waiting out the leftover interval // from before the reconnect. diff --git a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/WebSocketResponseTest.java b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/WebSocketResponseTest.java index b70bdc4f..b67b732b 100644 --- a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/WebSocketResponseTest.java +++ b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/WebSocketResponseTest.java @@ -39,7 +39,7 @@ public class WebSocketResponseTest { @Test public void testDurableAckFactory() throws Exception { assertMemoryLeak(() -> { - WebSocketResponse response = WebSocketResponse.durableAck("trades", 42L); + WebSocketResponse response = WebSocketResponse.durableAck("trades", "trades-dir", 42L); Assert.assertTrue(response.isDurableAck()); Assert.assertFalse(response.isSuccess()); Assert.assertEquals(1, response.getTableEntryCount()); @@ -54,7 +54,7 @@ public void testDurableAckFactory() throws Exception { @Test public void testDurableAckIsStructurallyValid() throws Exception { assertMemoryLeak(() -> { - WebSocketResponse response = WebSocketResponse.durableAck("t", 7L); + WebSocketResponse response = WebSocketResponse.durableAck("t", "t-dir", 7L); int size = response.serializedSize(); long ptr = Unsafe.malloc(size + 1, MemoryTag.NATIVE_DEFAULT); @@ -72,7 +72,7 @@ public void testDurableAckIsStructurallyValid() throws Exception { @Test public void testDurableAckRoundTripThroughNativeMemory() throws Exception { assertMemoryLeak(() -> { - WebSocketResponse original = WebSocketResponse.durableAck("orders", 12345L); + WebSocketResponse original = WebSocketResponse.durableAck("orders", "orders-dir", 12345L); int size = original.serializedSize(); long ptr = Unsafe.malloc(size, MemoryTag.NATIVE_DEFAULT); try { @@ -95,7 +95,7 @@ public void testDurableAckRoundTripThroughNativeMemory() throws Exception { @Test public void testDurableAckDoesNotCarryErrorMessage() throws Exception { assertMemoryLeak(() -> { - WebSocketResponse response = WebSocketResponse.durableAck("t", 99L); + WebSocketResponse response = WebSocketResponse.durableAck("t", "t-dir", 99L); int size = response.serializedSize(); long ptr = Unsafe.malloc(size, MemoryTag.NATIVE_DEFAULT); try {