diff --git a/CHANGES.txt b/CHANGES.txt index 4d74b0031..c3cda0a5a 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 559edc2f9..26c57279c 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 27c41ee11..e29882961 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.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 void testBuilderMethodCreatesValidBuilder() 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 574843241..630a1f3c5 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.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 static CdcBuilder builder(@NotNull String jobId, 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 cd8b3ea9e..f7c579660 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.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 CdcBuilder withCommitLogProvider(@NotNull CommitLogProvider commitLogProv return this; } - public CdcBuilder withStats(@NotNull CdcStats stats) + public CdcBuilder withStats(@NotNull ICdcStats stats) { this.stats = stats; return this;