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 dd4d2563 CASSANALYTICS-191: CDC reader stats silently dropped in 
SidecarCdcBuilder (#232)
dd4d2563 is described below

commit dd4d25637a50d9b20138e657972174b4bfb6fab3
Author: Jyothsna konisa <[email protected]>
AuthorDate: Mon Aug 10 12:19:39 2026 -0700

    CASSANALYTICS-191: CDC reader stats silently dropped in SidecarCdcBuilder 
(#232)
    
    Patch by Jyothsna Konisa; Reviewed by Saranya Krishnakumar for 
CASSANALYTICS-191
---
 CHANGES.txt                                        |  1 +
 .../cassandra/cdc/sidecar/SidecarCdcBuilder.java   |  1 +
 .../cassandra/cdc/sidecar/SidecarCdcTest.java      | 56 ++++++++++++++++++++++
 .../main/java/org/apache/cassandra/cdc/Cdc.java    | 12 +++++
 .../java/org/apache/cassandra/cdc/CdcBuilder.java  |  3 +-
 5 files changed, 71 insertions(+), 2 deletions(-)

diff --git a/CHANGES.txt b/CHANGES.txt
index 4d74b003..c3cda0a5 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
 0.5.0
 -----
+ * CDC reader stats silently dropped in SidecarCdcBuilder (CASSANALYTICS-191)
  * Add CapturePublishedSchema metric to SidecarCdcStats (CASSANALYTICS-189)
  * Expand list of architecture that supports unaligned access in 
FastByteOperations (CASSANALYTICS-188)
  * Fix FastByteOperations Silently Falling Back to Pure-Java Comparator 
(CASSANALYTICS-187)
diff --git 
a/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
 
b/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
index 559edc2f..26c57279 100644
--- 
a/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
+++ 
b/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
@@ -55,6 +55,7 @@ public class SidecarCdcBuilder extends CdcBuilder
         super(jobId, partitionId, eventConsumer, schemaSupplier);
         this.clusterConfigProvider = clusterConfigProvider;
         this.sidecarCdcClient = sidecarCdcClient;
+        withStats(cdcStats);
         withCdcOptions(cdcOptions);
         withTokenRangeSupplier(tokenRangeSupplier);
     }
diff --git 
a/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
 
b/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
index 27c41ee1..e2988296 100644
--- 
a/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
+++ 
b/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
@@ -20,6 +20,8 @@
 package org.apache.cassandra.cdc.sidecar;
 
 import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.CompletableFuture;
 
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
@@ -31,10 +33,14 @@ import org.apache.cassandra.cdc.api.EventConsumer;
 import org.apache.cassandra.cdc.api.SchemaSupplier;
 import org.apache.cassandra.cdc.api.TokenRangeSupplier;
 import org.apache.cassandra.cdc.stats.ICdcStats;
+import org.apache.cassandra.spark.data.CqlTable;
+import org.apache.cassandra.spark.data.ReplicationFactor;
 import org.apache.cassandra.spark.data.partitioner.CassandraInstance;
+import org.apache.cassandra.spark.utils.AsyncExecutor;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
 
 /**
  * Unit tests for SidecarCdc class
@@ -81,6 +87,56 @@ public class SidecarCdcTest
         assertThat(builder.sidecarCdcClient).isEqualTo(mockSidecarCdcClient);
     }
 
+    /**
+     * Regression test for a bug where {@link SidecarCdcBuilder}'s constructor 
accepted an
+     * {@link ICdcStats} parameter but never wired it into the builder (via 
{@link SidecarCdcBuilder#withStats}),
+     * so every {@link SidecarCdc} built through it silently used {@link 
ICdcStats#STUB} instead of the real,
+     * caller-supplied stats implementation — with no exception anywhere to 
reveal it. This mirrors the exact
+     * call shape a consuming project (e.g. cassandra-sidecar's CdcManager) 
uses in production:
+     * {@code 
SidecarCdc.builder(...).withExecutor(...).withReplicationFactorSupplier(...).withSidecarStatePersister(...).build()}.
+     */
+    @Test
+    public void testBuiltSidecarCdcUsesSuppliedStatsNotStub() throws Exception
+    {
+        String jobId = "test-job-123";
+        int partitionId = 0;
+        CdcOptions cdcOptions = mock(CdcOptions.class);
+        ClusterConfigProvider clusterConfigProvider = 
mock(ClusterConfigProvider.class);
+        when(clusterConfigProvider.dc()).thenReturn("DC1");
+        EventConsumer eventConsumer = mock(EventConsumer.class);
+        TokenRangeSupplier tokenRangeSupplier = mock(TokenRangeSupplier.class);
+        SidecarCdcClient mockSidecarCdcClient = mock(SidecarCdcClient.class);
+        AsyncExecutor asyncExecutor = mock(AsyncExecutor.class);
+
+        // Just enough of a CDC-enabled table (with a replication factor for 
"DC1") to satisfy
+        // SidecarCdc.initSchema(), which runs synchronously inside the 
constructor.
+        ReplicationFactor rf = new 
ReplicationFactor(ReplicationFactor.ReplicationStrategy.NetworkTopologyStrategy,
+                                                     Map.of("DC1", 3));
+        CqlTable cqlTable = mock(CqlTable.class);
+        when(cqlTable.replicationFactor()).thenReturn(rf);
+        SchemaSupplier schemaSupplier = mock(SchemaSupplier.class);
+        
when(schemaSupplier.getCDCEnabledTables()).thenReturn(CompletableFuture.completedFuture(Set.of(cqlTable)));
+
+        SidecarCdc consumer = SidecarCdc.builder(jobId,
+                                                 partitionId,
+                                                 cdcOptions,
+                                                 clusterConfigProvider,
+                                                 eventConsumer,
+                                                 schemaSupplier,
+                                                 tokenRangeSupplier,
+                                                 mockSidecarCdcClient,
+                                                 cdcStats)
+                                        .withExecutor(asyncExecutor)
+                                        .build();
+
+        assertThat(consumer.stats())
+            .as("SidecarCdc.builder(...)'s cdcStats argument must reach 
Cdc.stats — if it doesn't, every "
+                + "ICdcStats call (changeProduced, insufficientReplicas, 
mutationsReadCount, etc.) silently "
+                + "no-ops against ICdcStats.STUB instead of the real 
implementation, with no exception to reveal it.")
+            .isSameAs(cdcStats)
+            .isNotSameAs(ICdcStats.STUB);
+    }
+
     @Test
     public void testPerInstancePortResolution()
     {
diff --git 
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java 
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
index 57484324..630a1f3c 100644
--- a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
+++ b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
@@ -36,6 +36,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import com.esotericsoftware.kryo.io.Output;
+import com.google.common.annotations.VisibleForTesting;
 import org.apache.cassandra.bridge.CassandraBridge;
 import org.apache.cassandra.bridge.CdcBridge;
 import org.apache.cassandra.bridge.CdcBridgeFactory;
@@ -110,6 +111,17 @@ public class Cdc implements Closeable
         return new CdcBuilder(jobId, partitionId, eventConsumer, 
schemaSupplier);
     }
 
+    /**
+     * @return the {@link ICdcStats} this {@link Cdc} instance was built with. 
Exposed so tests can assert
+     * the stats implementation supplied to the builder is actually the 
instance in use, and not the
+     * {@link ICdcStats#STUB} default silently falling through.
+     */
+    @VisibleForTesting
+    public ICdcStats stats()
+    {
+        return stats;
+    }
+
     public String jobId()
     {
         return jobId;
diff --git 
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
 
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
index cd8b3ea9..f7c57966 100644
--- 
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
+++ 
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
@@ -31,7 +31,6 @@ import org.apache.cassandra.cdc.api.SchemaSupplier;
 import org.apache.cassandra.cdc.api.StatePersister;
 import org.apache.cassandra.cdc.api.TableIdLookup;
 import org.apache.cassandra.cdc.api.TokenRangeSupplier;
-import org.apache.cassandra.cdc.stats.CdcStats;
 import org.apache.cassandra.cdc.stats.ICdcStats;
 import org.apache.cassandra.spark.utils.AsyncExecutor;
 import org.jetbrains.annotations.NotNull;
@@ -131,7 +130,7 @@ public class CdcBuilder
         return this;
     }
 
-    public CdcBuilder withStats(@NotNull CdcStats stats)
+    public CdcBuilder withStats(@NotNull ICdcStats stats)
     {
         this.stats = stats;
         return this;


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

Reply via email to