Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,7 @@ public class WebSocketResponse {
// subsequent sight is allocation-free.
private final Utf8SequenceObjHashMap<String> tableNameCache = new Utf8SequenceObjHashMap<>();
private final ObjList<String> tableNames = new ObjList<>();
private final ObjList<String> tableDirNames = new ObjList<>();
private final LongList tableSeqTxns = new LongList();
private String errorMessage;
private int errorMessageUtf8Length;
Expand All @@ -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;
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -272,6 +288,7 @@ public boolean isSuccess() {
*/
public boolean readFrom(long ptr, int length) {
tableNames.clear();
tableDirNames.clear();
tableSeqTxns.clear();

if (length < 1) {
Expand All @@ -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
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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.
* <p>
* 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);
Expand Down Expand Up @@ -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();
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, String> 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.
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand All @@ -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);
Expand All @@ -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 {
Expand All @@ -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 {
Expand Down
Loading