This is an automated email from the ASF dual-hosted git repository.

jyothsnakonisa pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra-analytics.git


The following commit(s) were added to refs/heads/trunk by this push:
     new c0aab1d4 CASSANALYTICS-184 : CdcState.ReplicaCountSerializer map-size 
overflow corrupts persisted CDC state (#225)
c0aab1d4 is described below

commit c0aab1d44a510fd26527779747d22e8b378d00e6
Author: Jyothsna konisa <[email protected]>
AuthorDate: Tue Jul 21 07:03:59 2026 -0700

    CASSANALYTICS-184 : CdcState.ReplicaCountSerializer map-size overflow 
corrupts persisted CDC state (#225)
    
    Patch by Jyothsna Konisa; Reviewed by Josh McKenzie for CASSANALYTICS-184
---
 CHANGES.txt                                        |  1 +
 .../cdc/model/CdcKryoSerializationTests.java       | 46 ++++++++++++++++++++++
 .../org/apache/cassandra/cdc/state/CdcState.java   | 22 +++++++++--
 3 files changed, 66 insertions(+), 3 deletions(-)

diff --git a/CHANGES.txt b/CHANGES.txt
index 49ed57ae..811417a7 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
 0.5.0
 -----
+ * CdcState.ReplicaCountSerializer map-size overflow corrupts persisted CDC 
state (CASSANALYTICS-184)
  * SSTable-version-based bridge determination (CASSANALYTICS-24)
  * Upgrade sidecar version to 0.4.0 (CASSANALYTICS-176)
  * Exclude IP address from RingInstance equality so node replacement does not 
fail bulk write jobs (CASSANALYTICS-175)
diff --git 
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/model/CdcKryoSerializationTests.java
 
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/model/CdcKryoSerializationTests.java
index 66c076cb..58c1bf7b 100644
--- 
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/model/CdcKryoSerializationTests.java
+++ 
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/model/CdcKryoSerializationTests.java
@@ -20,6 +20,8 @@
 package org.apache.cassandra.cdc.model;
 
 import java.math.BigInteger;
+import java.util.HashMap;
+import java.util.Map;
 
 import com.google.common.collect.ImmutableMap;
 import org.junit.jupiter.api.Test;
@@ -110,4 +112,48 @@ public class CdcKryoSerializationTests
         assertThat(deserialized).isNotNull();
         assertThat(deserialized).isEqualTo(expected);
     }
+
+    /**
+     * Regression test for the {@code ReplicaCountSerializer} short-overflow 
bug: the replica-count
+     * map size used to be written/read as a signed short (max 32767), so a 
watermarker larger than
+     * that (reachable via {@code CdcOptions.maxCdcStateSize()}/{@code 
maxWatermarkerSize()}, both
+     * configurable well above 32767) would silently overflow on write and 
deserialize with a
+     * corrupted (potentially negative) size, observed in production as a 
permanent
+     * restart-crash-loop ({@code IllegalArgumentException: Illegal initial 
capacity}). The map size
+     * is now a signed int, so this must round-trip correctly for a map with 
more than 32767 entries.
+     */
+    @Test
+    public void testCdcStateReplicaCountMapLargerThanShortOverflows()
+    {
+        int size = Short.MAX_VALUE + 1000; // comfortably past the old 32767 
short limit
+        Map<PartitionUpdateWrapper.Digest, Integer> replicaCount = new 
HashMap<>(size);
+        for (int i = 0; i < size; i++)
+        {
+            PartitionUpdateWrapper.Digest digest = new 
PartitionUpdateWrapper.Digest("ks1",
+                                                                               
       "tb1",
+                                                                               
       TimeUtils.nowMicros(),
+                                                                               
       intToBytes(i),
+                                                                               
       500,
+                                                                               
       BigInteger.ONE);
+            replicaCount.put(digest, i % 128); // replica count is serialized 
as a single byte
+        }
+
+        CdcState expected = CdcState.of(1L,
+                                        TokenRange.openClosed(BigInteger.ONE, 
BigInteger.TEN),
+                                        CommitLogMarkers.EMPTY,
+                                        replicaCount);
+        assertThat(expected.size()).isEqualTo(size);
+
+        testCdcStateSerialization(expected);
+    }
+
+    private static byte[] intToBytes(int value)
+    {
+        return new byte[]{
+        (byte) (value >>> 24),
+        (byte) (value >>> 16),
+        (byte) (value >>> 8),
+        (byte) value
+        };
+    }
 }
diff --git 
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/cdc/state/CdcState.java
 
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/cdc/state/CdcState.java
index 2d48ace1..c8a15cb5 100644
--- 
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/cdc/state/CdcState.java
+++ 
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/cdc/state/CdcState.java
@@ -395,12 +395,28 @@ public class CdcState
         }
     }
 
+    /**
+     * Serializes the replica-count map for late/un-acked mutations awaiting 
sufficient replicas.
+     * Map key = mutation digest; value = replica count seen so far for that 
mutation (always small,
+     * bounded by replication factor, so kept as a single byte -- unrelated to 
the bug below).
+     *
+     * <p><b>Wire format note:</b> the map <em>size</em> (entry count, i.e. 
number of distinct
+     * mutations tracked) is written as a signed int -- previously a signed 
short (max 32767
+     * entries), which silently overflowed and corrupted persisted state once
+     * {@code CdcOptions.maxCdcStateSize()}/{@code maxWatermarkerSize()} 
configured a larger value
+     * (both default well above 32767), observed in production as a permanent 
restart-crash-loop.
+     * {@code int} is sufficient since the configured max size is always far 
below
+     * {@code Integer.MAX_VALUE}; {@code long} would be unwarranted here. This 
is a wire-format
+     * change: state persisted by the old (short) format cannot be read by 
this version and vice
+     * versa. Rolling upgrades should account for this (e.g. purge/rebuild 
persisted CDC state
+     * rather than carry it across the upgrade boundary).
+     */
     public static class ReplicaCountSerializer extends 
com.esotericsoftware.kryo.Serializer<Map<PartitionUpdateWrapper.Digest, 
Integer>>
     {
         public Map<PartitionUpdateWrapper.Digest, Integer> read(Kryo kryo, 
Input in, Class type)
         {
-            // read replica counts
-            int numUpdates = in.readShort();
+            // numUpdates = number of distinct mutation digests tracked (map 
size), not a replica count
+            int numUpdates = in.readInt();
             Map<PartitionUpdateWrapper.Digest, Integer> replicaCounts = new 
HashMap<>(numUpdates);
             for (int i = 0; i < numUpdates; i++)
             {
@@ -413,7 +429,7 @@ public class CdcState
         public void write(Kryo kryo, Output out, 
Map<PartitionUpdateWrapper.Digest, Integer> o)
         {
             // write replica counts for late mutations
-            out.writeShort(o.size());
+            out.writeInt(o.size());
             for (Map.Entry<PartitionUpdateWrapper.Digest, Integer> entry : 
o.entrySet())
             {
                 PartitionUpdateWrapper.Digest digest = entry.getKey();


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to