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-sidecar.git


The following commit(s) were added to refs/heads/trunk by this push:
     new df567616 CASSSIDECAR-483: CDC configs in configs table are not wired 
to SidecarCdcOptions/SidecarStatePersister (#370)
df567616 is described below

commit df567616145b426296e29331f6e065013d5c956c
Author: Jyothsna konisa <[email protected]>
AuthorDate: Tue Jul 21 11:36:57 2026 -0700

    CASSSIDECAR-483: CDC configs in configs table are not wired to 
SidecarCdcOptions/SidecarStatePersister (#370)
    
    Patch by Jyothsna Konisa; Reviewed by Josh McKenzie for CASSSIDECAR-483
---
 CHANGES.txt                                        |  1 +
 .../apache/cassandra/sidecar/cdc/CdcManager.java   | 37 ++++++++-
 .../cassandra/sidecar/cdc/SidecarCdcOptions.java   | 64 ++++++++++++++-
 .../cassandra/sidecar/modules/CdcModule.java       |  4 +-
 .../cassandra/sidecar/cdc/CdcManagerTest.java      | 22 +++++
 .../sidecar/cdc/SidecarCdcOptionsTest.java         | 93 ++++++++++++++++++++++
 6 files changed, 215 insertions(+), 6 deletions(-)

diff --git a/CHANGES.txt b/CHANGES.txt
index eda42dee..d9caa34b 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
 0.5.0
 -----
+ * Wire CDC configs in configs table to 
SidecarCdcOptions/SidecarStatePersister (CASSSIDECAR-483)
  * Implement durable operational job tracker (CASSSIDECAR-374)
  * Remove filesystem path from Http response (CASSSIDECAR-477)
  * Add basic configuration retrieval logic to ConfigurationManager 
(CASSSIDECAR-427)
diff --git 
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java 
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
index 88f05d02..8c97b1fd 100644
--- a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
+++ b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
@@ -18,6 +18,7 @@
 
 package org.apache.cassandra.sidecar.cdc;
 
+import java.time.Duration;
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
@@ -35,6 +36,7 @@ import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
 import org.apache.cassandra.cdc.sidecar.ReplicationFactorSupplier;
 import org.apache.cassandra.cdc.sidecar.SidecarCdc;
 import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcOptions;
 import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
 import org.apache.cassandra.cdc.sidecar.SidecarStatePersister;
 import org.apache.cassandra.cdc.stats.ICdcStats;
@@ -219,10 +221,43 @@ public class CdcManager
 
     private @NotNull SidecarStatePersister getSidecarStatePersister()
     {
-        return new 
SidecarStatePersister(org.apache.cassandra.cdc.sidecar.SidecarCdcOptions.DEFAULT,
+        return new SidecarStatePersister(new 
ConfigBackedPersisterOptions(conf),
                                          cdcOptions,
                                          SidecarCdcStats.STUB,
                                          cassandraClient,
                                          asyncExecutor);
     }
+
+    /**
+     * Adapts the DB-backed {@link CdcConfig} to the 
cassandra-analytics-cdc-sidecar
+     * {@link SidecarCdcOptions} interface consumed by {@link 
SidecarStatePersister}, which only
+     * ever calls {@link SidecarCdcOptions#persistDelay()} on it.
+     *
+     * <p>Previously {@link SidecarStatePersister} was built with {@code 
SidecarCdcOptions.DEFAULT},
+     * which pinned {@code persistDelay()} to its hardcoded 1000ms interface 
default regardless of
+     * what operators configured in the "configs" table.
+     *
+     * <p>This is a separate, minimal implementation rather than reusing this 
class's own
+     * {@code SidecarCdcOptions} (the {@link CdcOptions} implementation 
consumed by the CDC read
+     * path): the two interfaces have different call sites and no requirement 
to be backed by the
+     * same object, so keeping them independent avoids one implementation 
growing overrides it
+     * doesn't need to satisfy the other's contract.
+     */
+    @VisibleForTesting
+    static final class ConfigBackedPersisterOptions implements 
SidecarCdcOptions
+    {
+        private final CdcConfig conf;
+
+        @VisibleForTesting
+        ConfigBackedPersisterOptions(CdcConfig conf)
+        {
+            this.conf = conf;
+        }
+
+        @Override
+        public Duration persistDelay()
+        {
+            return Duration.ofMillis(conf.persistDelay().toMillis());
+        }
+    }
 }
diff --git 
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptions.java 
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptions.java
index 10dcac16..0b9c8a21 100644
--- 
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptions.java
+++ 
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptions.java
@@ -18,6 +18,7 @@
 
 package org.apache.cassandra.sidecar.cdc;
 
+import java.time.Duration;
 import java.util.Map;
 
 import org.apache.cassandra.bridge.CassandraVersion;
@@ -26,16 +27,22 @@ import 
org.apache.cassandra.sidecar.utils.InstanceMetadataFetcher;
 import org.apache.cassandra.spark.data.ReplicationFactor;
 
 /**
- * Specific sidecar CDC options
+ * Specific sidecar CDC options, consumed by the CDC read path ({@code 
SidecarCdc}).
+ *
+ * <p>Delegates the throughput/backpressure-related knobs to {@link CdcConfig} 
so that they are
+ * live-tunable from the DB-backed "configs" table (via {@link 
org.apache.cassandra.sidecar.tasks.CdcConfigRefresherNotifierTask})
+ * without a Sidecar restart, instead of silently falling back to the {@link 
CdcOptions} interface
+ * defaults baked into the cassandra-analytics library.
  */
 public class SidecarCdcOptions implements CdcOptions
 {
-
     private final InstanceMetadataFetcher instanceMetadataFetcher;
+    private final CdcConfig conf;
 
-    public SidecarCdcOptions(InstanceMetadataFetcher instanceMetadataFetcher)
+    public SidecarCdcOptions(InstanceMetadataFetcher instanceMetadataFetcher, 
CdcConfig conf)
     {
         this.instanceMetadataFetcher = instanceMetadataFetcher;
+        this.conf = conf;
     }
 
 
@@ -59,4 +66,55 @@ public class SidecarCdcOptions implements CdcOptions
                 instance -> 
instance.delegate().nodeSettings().releaseVersion());
         return 
CassandraVersion.fromVersion(releaseVersion).orElse(CassandraVersion.FOURZERO);
     }
+
+    /**
+     * Add an optional delay between micro-batches, to slow CDC down if it is 
overwhelming Cassandra
+     * or the downstream Kafka publish stage. Backed by {@code 
CdcConfig.minDelayBetweenMicroBatches()}
+     * so it can be lowered/raised live via the "configs" table, e.g. to 
accelerate cdc_raw drain
+     * during a backlog without a restart.
+     */
+    @Override
+    public Duration minimumDelayBetweenMicroBatches()
+    {
+        return 
Duration.ofMillis(conf.minDelayBetweenMicroBatches().toMillis());
+    }
+
+    /**
+     * Throttles how many commit logs are read per epoch per instance. Backed 
by
+     * {@code CdcConfig.maxCommitLogsPerInstance()} so it can be raised live 
to catch up faster
+     * on a backlog, or lowered to bound per-batch memory/duration during 
burst load.
+     */
+    @Override
+    public int maxCommitLogsPerInstance()
+    {
+        return conf.maxCommitLogsPerInstance();
+    }
+
+    /**
+     * Maximum number of late/un-acked mutation digests held in the CDC 
watermarker state. Backed
+     * by {@code CdcConfig.maxWatermarkerSize()}.
+     *
+     * <p><b>Caution:</b> {@code CdcState.ReplicaCountSerializer} currently 
serializes this map's
+     * size with {@code writeShort}/{@code readShort} (signed 16-bit, max 
32767). Do not configure
+     * this above 32767 until that serializer is widened to an int, or 
persisted CDC state can
+     * silently corrupt (observed as a permanent restart-crash-loop in 
production).
+     */
+    @Override
+    public int maxCdcStateSize()
+    {
+        return conf.maxWatermarkerSize();
+    }
+
+    /**
+     * Maximum age of mutations retained in the CDC watermarker before being 
purged (and counted via
+     * {@code droppedExpiredMutations}). Backed by {@code 
CdcConfig.watermarkWindow()} -- previously
+     * this value was entirely unreachable: {@code watermarkWindow()} was read 
from the DB-backed
+     * config but never consulted by the CDC engine, which instead silently 
used the 1-hour
+     * {@link CdcOptions#maximumAge()} interface default regardless of what 
operators configured.
+     */
+    @Override
+    public Duration maximumAge()
+    {
+        return Duration.ofSeconds(conf.watermarkWindow().toSeconds());
+    }
 }
diff --git 
a/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java 
b/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
index e49538e1..9d1811d4 100644
--- a/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
+++ b/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
@@ -490,9 +490,9 @@ public class CdcModule extends AbstractModule
 
     @Provides
     @Singleton
-    public CdcOptions cdcOptions(InstanceMetadataFetcher 
instanceMetadataFetcher)
+    public CdcOptions cdcOptions(InstanceMetadataFetcher 
instanceMetadataFetcher, CdcConfig conf)
     {
-        return new SidecarCdcOptions(instanceMetadataFetcher);
+        return new SidecarCdcOptions(instanceMetadataFetcher, conf);
     }
 
     @ProvidesIntoMap
diff --git 
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java 
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
index 55d7c169..0701e855 100644
--- a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
+++ b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
@@ -22,12 +22,14 @@ import java.io.IOException;
 import java.math.BigInteger;
 import java.net.InetAddress;
 import java.net.UnknownHostException;
+import java.time.Duration;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.TimeUnit;
 
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
@@ -42,6 +44,7 @@ import org.apache.cassandra.cdc.sidecar.SidecarStatePersister;
 import org.apache.cassandra.cdc.stats.ICdcStats;
 import org.apache.cassandra.sidecar.cluster.instance.InstanceMetadata;
 import org.apache.cassandra.sidecar.common.server.cluster.locator.TokenRange;
+import 
org.apache.cassandra.sidecar.common.server.utils.MillisecondBoundConfiguration;
 import org.apache.cassandra.sidecar.concurrent.TaskExecutorPool;
 import org.apache.cassandra.sidecar.coordination.RangeManager;
 import org.apache.cassandra.sidecar.db.CdcDatabaseAccessor;
@@ -369,6 +372,25 @@ public class CdcManagerTest
         assertThat(cdcManager.getInstanceId(unknownIp)).isEqualTo(-1);
     }
 
+    /**
+     * Regression guard: {@code SidecarStatePersister} was previously built 
with the
+     * cassandra-analytics-cdc-sidecar {@code SidecarCdcOptions.DEFAULT}, 
which pinned
+     * {@code persistDelay()} to its hardcoded 1000ms interface default 
regardless of what
+     * operators configured in the "configs" table. {@link 
CdcManager.ConfigBackedPersisterOptions}
+     * fixes this by delegating {@code persistDelay()} straight to {@link 
CdcConfig}; this test
+     * uses a value that differs from both the interface default (1000ms) and 
the
+     * {@code CdcConfigImpl} default (also 1000ms) so a pass proves real 
delegation.
+     */
+    @Test
+    void configBackedPersisterOptionsDelegatesPersistDelayToCdcConfig()
+    {
+        when(cdcConfig.persistDelay()).thenReturn(new 
MillisecondBoundConfiguration(2500, TimeUnit.MILLISECONDS));
+
+        CdcManager.ConfigBackedPersisterOptions persisterOptions = new 
CdcManager.ConfigBackedPersisterOptions(cdcConfig);
+
+        
assertThat(persisterOptions.persistDelay()).isEqualTo(Duration.ofMillis(2500));
+    }
+
     // Helper methods
 
     private TokenRange mockTokenRange(BigInteger start, BigInteger end)
diff --git 
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptionsTest.java
 
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptionsTest.java
new file mode 100644
index 00000000..29fd6d34
--- /dev/null
+++ 
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/SidecarCdcOptionsTest.java
@@ -0,0 +1,93 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.cassandra.sidecar.cdc;
+
+import java.time.Duration;
+import java.util.concurrent.TimeUnit;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import 
org.apache.cassandra.sidecar.common.server.utils.MillisecondBoundConfiguration;
+import 
org.apache.cassandra.sidecar.common.server.utils.SecondBoundConfiguration;
+import org.apache.cassandra.sidecar.utils.InstanceMetadataFetcher;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * Unit tests for {@link SidecarCdcOptions}.
+ *
+ * <p>These specifically guard against the throughput/backpressure knobs 
silently falling back to
+ * the hardcoded {@code CdcOptions} interface defaults (1000ms micro-batch 
delay, 4 commit logs
+ * per instance, 200000 max state size, 1 hour max age) instead of the 
DB-backed {@link CdcConfig}
+ * value an operator configures in the "configs" table. Every value asserted 
here is deliberately
+ * chosen to differ from both the {@code CdcOptions} interface default and the 
{@link CdcConfig}
+ * default, so a passing test proves real delegation rather than two defaults 
coincidentally
+ * matching.
+ */
+class SidecarCdcOptionsTest
+{
+    private CdcConfig conf;
+    private SidecarCdcOptions options;
+
+    @BeforeEach
+    void setUp()
+    {
+        conf = mock(CdcConfig.class);
+        InstanceMetadataFetcher instanceMetadataFetcher = 
mock(InstanceMetadataFetcher.class);
+        options = new SidecarCdcOptions(instanceMetadataFetcher, conf);
+    }
+
+    @Test
+    void minimumDelayBetweenMicroBatchesDelegatesToCdcConfig()
+    {
+        when(conf.minDelayBetweenMicroBatches()).thenReturn(new 
MillisecondBoundConfiguration(250, TimeUnit.MILLISECONDS));
+
+        
assertThat(options.minimumDelayBetweenMicroBatches()).isEqualTo(Duration.ofMillis(250));
+    }
+
+    @Test
+    void maxCommitLogsPerInstanceDelegatesToCdcConfig()
+    {
+        when(conf.maxCommitLogsPerInstance()).thenReturn(16);
+
+        assertThat(options.maxCommitLogsPerInstance()).isEqualTo(16);
+    }
+
+    @Test
+    void maxCdcStateSizeDelegatesToCdcConfigMaxWatermarkerSize()
+    {
+        when(conf.maxWatermarkerSize()).thenReturn(12345);
+
+        assertThat(options.maxCdcStateSize()).isEqualTo(12345);
+    }
+
+    @Test
+    void maximumAgeDelegatesToCdcConfigWatermarkWindow()
+    {
+        // Regression guard: watermarkWindow() was previously read from the 
DB-backed config but
+        // never consulted anywhere, so operators configuring it had no actual 
effect -- the CDC
+        // engine silently used the 1-hour CdcOptions interface default 
instead.
+        when(conf.watermarkWindow()).thenReturn(new 
SecondBoundConfiguration(120, TimeUnit.SECONDS));
+
+        assertThat(options.maximumAge()).isEqualTo(Duration.ofSeconds(120));
+    }
+}


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

Reply via email to