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
18 changes: 14 additions & 4 deletions .github/workflows/test.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ jobs:
path: ${{ github.workspace }}
key: build-jdk11-${{ github.sha }}

# JDK17 dependency build for the Spark 4 / Scala 2.13 / Cassandra 5.0 matrix.
# JDK17 dependency build for the Spark 4 / Scala 2.13 / Cassandra 5.0 and 6.0 matrix.
# Produces a distinct workspace cache (build-jdk17-...) so downstream test jobs
# must restore from exactly one of build-jdk11 / build-jdk17 — mixing them
# would clobber dependency jars.
Expand Down Expand Up @@ -117,7 +117,7 @@ jobs:
export SCALA_VERSION="2.13"
# JDK17 only targets Cassandra 5.0+; skip 4.0 / 4.1 dtest jar builds
# (build-dtest-jars.sh reads this var to filter CANDIDATE_BRANCHES).
export BRANCHES="cassandra-5.0"
export BRANCHES="cassandra-5.0 cassandra-6.0"

./scripts/build-dependencies.sh

Expand Down Expand Up @@ -217,16 +217,21 @@ jobs:
needs: [build-jdk11, build-jdk17]
runs-on: ubuntu-latest
strategy:
# GitHub Actions generate a cross-product of 'config' × 'job_index' (4 × 5 = 20 jobs).
# GitHub Actions generate a cross-product of 'config' × 'job_index' (6 × 5 = 30 jobs).
# The 'include' entries don't add new combinations — they augment existing ones
# by matching on 'config' and injecting 'scala', 'cassandra', 'jdk', and 'spark'
# into each match. To add a new version: add one entry to 'config' and one to
# 'include'.
matrix:
config: ['s2.13-c5.0.7', 's2.12-c4.1.4', 's2.12-c4.0.17', 's2.13-c5.0.7-spark4']
config: ['s2.13-c6.0a2', 's2.13-c5.0.7', 's2.12-c4.1.4', 's2.12-c4.0.17', 's2.13-c6.0a2-spark4', 's2.13-c5.0.7-spark4']
job_index: [0, 1, 2, 3, 4]
job_total: [5]
include:
- config: 's2.13-c6.0a2'
scala: '2.13'
cassandra: '6.0-alpha2'
jdk: '11'
spark: '3'
- config: 's2.13-c5.0.7'
scala: '2.13'
cassandra: '5.0.7'
Expand All @@ -242,6 +247,11 @@ jobs:
cassandra: '4.0.17'
jdk: '11'
spark: '3'
- config: 's2.13-c6.0a2-spark4'
scala: '2.13'
cassandra: '6.0-alpha2'
jdk: '17'
spark: '4'
- config: 's2.13-c5.0.7-spark4'
scala: '2.13'
cassandra: '5.0.7'
Expand Down
26 changes: 13 additions & 13 deletions DEV-README.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,21 +125,21 @@ Without these entries an upstream DNS server can answer `localhost2` with a publ
connects to that address and the test fails with `java.net.ConnectException: Operation timed out` in
`CassandraDataLayer.initialize`.

### Topology-change tests skip on Cassandra 6.0
### Topology-change tests

Twenty-eight test classes under `expansion`, `shrink`, `replacement` and `movement` pause a topology change with a
ByteBuddy hook, then run the bulk writer while the node is in the transitional state. The hooks target
`StorageService.bootstrap(Collection, long)`, `StorageService.unbootstrap()` and `RangeRelocator.stream()`.
The tests under `expansion`, `shrink`, `replacement` and `movement` pause a topology change with ByteBuddy, then run the bulk writer while the node is in the transitional state. `TopologyChangeBBUtils` selects the hook from the classes available in the Cassandra instance classloader and fails immediately if the target method is absent.

CEP-21 Transactional Cluster Metadata removed all three in Cassandra 6.0. The work now belongs to
`org.apache.cassandra.tcm.sequences`: `BootstrapAndJoin.bootstrap(...)` for a join, `BootstrapAndReplace` for a
replacement, `UnbootstrapAndLeave.executeNext()` with `LeaveStreams` for a decommission, and `Move` for a token move.
A hook that fails to install is silent, so each class waited two minutes for a latch that never counted down.
| Operation | Cassandra 4.x / 5.0 | Cassandra 6.0 (Transactional Cluster Metadata) |
| --- | --- | --- |
| Join or replace | `StorageService.bootstrap(Collection, long)` | `BootstrapAndJoin.bootstrap(...)`, also called by replacement |
| Decommission | `StorageService.unbootstrap()` | `UnbootstrapStreams.execute(...)` |
| Move | `RangeRelocator.stream()` | `Move.executeNext()` |

`ResiliencyTestBase.assumeTopologyChangeHooksSupported()` now skips these classes on 6.0 and later. The four base
classes call it from `beforeClusterProvisioning()`. Bulk write during a topology change is therefore untested on
6.0. To close the gap, retarget each hook at the sequence types named above, and keep the 4.0 and 5.0 targets for
the older runs.
The shared single-node scenarios create their schema before starting a topology change. On Cassandra 6.0 and later, they configure three Cluster Metadata Service members before a node can leave or stop. These scenarios also cover multi-datacenter clusters.

The `JoiningDisjointRangesTest` and `LeavingDisjointRangesTest` scenarios exercise two concurrent transitions on Cassandra 6.0 and later. They use Sidecar's CASSSIDECAR-277 token placement: six tokens per datacenter with the tokens for nodes 6 and 11 exchanged. The single-replicated variant also checks that the joining node in the unreplicated datacenter receives no rows. Successful joins must reach `Normal`; decommission tests await the nodetool result and check success or the injected failure before validating data again.

Cassandra versions older than 6.0 additionally run the concurrent join/leave and cluster-doubling/halving tests. The `LegacyJoiningMultiDC*` and `LegacyLeavingMultiDC*` tests preserve the original concurrent multi-datacenter scenarios alongside the shared single-node tests. These legacy scenarios retain schema creation during the transition and enable their required range-movement overrides only for the lifetime of the test class. They skip before cluster provisioning on Cassandra 6.0 and later because their concurrent transitions can require overlapping range locks.

## IntelliJ

Expand Down Expand Up @@ -172,4 +172,4 @@ compiled classes of both earlier majors, and Gradle never recompiles an inherite
changed between two majors is silent at build time and fails at runtime with `NoSuchMethodError`. After you add a
new major version, compile the union of the inherited sources against the new shaded jar and confirm that the
resulting error set matches the error set that the previous major produces. Every error the new version adds names
a file you must copy forward and override.
a file you must copy forward and override.
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,7 @@ public AbstractCluster<I> initializeCluster(String versionString,
{
instanceConfigUpdater = instanceConfigUpdater.andThen(config -> configuration.additionalInstanceConfig.forEach(config::set));
}
clusterBuilder.withTokenSupplier(tokenSupplier)
clusterBuilder.withTokenSupplier(configuration.tokenSupplier == null ? tokenSupplier : configuration.tokenSupplier)
.withConfig(instanceConfigUpdater);

if (dcCount > 1)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -190,6 +190,7 @@ protected void setup() throws Exception
assertThat(cluster).isNotNull();
afterClusterProvisioned();
initializeSchemaForTest();
afterSchemaInitialized();
mtlsTestHelper = new MtlsTestHelper(secretsPath);
startSidecar(cluster);
beforeTestStart();
Expand Down Expand Up @@ -309,6 +310,13 @@ protected void afterClusterProvisioned()
*/
protected abstract void initializeSchemaForTest();

/**
* Runs after schema creation, before Sidecar starts.
*/
protected void afterSchemaInitialized()
{
}

/**
* Override to perform an action before the tests start
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import com.google.common.base.Preconditions;

import org.apache.cassandra.distributed.api.Feature;
import org.apache.cassandra.distributed.api.TokenSupplier;
import org.apache.cassandra.distributed.shared.NetworkTopology;

/**
Expand All @@ -44,6 +45,7 @@ public class ClusterBuilderConfiguration
public String partitioner;
public Map<String, Object> additionalInstanceConfig = null;
public int tokenCount = 1;
public TokenSupplier tokenSupplier;
public IntFunction<NetworkTopology.DcAndRack> dcAndRackSupplier;

/**
Expand Down Expand Up @@ -180,6 +182,18 @@ public ClusterBuilderConfiguration tokenCount(int tokenCount)
return this;
}

/**
* Overrides the default token allocation for topology-change scenarios.
*
* @param tokenSupplier the token allocation for existing and joining nodes
* @return this configuration instance
*/
public ClusterBuilderConfiguration tokenSupplier(TokenSupplier tokenSupplier)
{
this.tokenSupplier = tokenSupplier;
return this;
}

/**
* Sets a supplier function that provides datacenter and rack information for each node in the cluster.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,8 +107,6 @@ public static void configureDefaultDTestJarProperties()
{
// Settings to reduce the test setup delay incurred if gossip is enabled
System.setProperty("cassandra.ring_delay_ms", "5000"); // down from 30s default
System.setProperty("cassandra.consistent.rangemovement", "false");
System.setProperty("cassandra.consistent.simultaneousmoves.allow", "true");
// End gossip delay settings
// Set the location of dtest jars
System.setProperty("cassandra.test.dtest_jar_path", System.getProperty("cassandra.test.dtest_jar_path", "dependencies"));
Expand All @@ -125,11 +123,6 @@ public static void configureDefaultDTestJarProperties()
System.setProperty("cassandra.require_native_file_hints", "true");
// Disable all native stuff in Netty as streaming isn't functional with native enabled
System.setProperty("shaded.io.netty.transport.noNative", "true");
// Lifted from the Simulation runner (we're running into similar errors):
// this property is used to allow non-members of the ring to exist in gossip without breaking RF changes
// it would be nice not to rely on this, but hopefully we'll have consistent range movements before it matters
System.setProperty("cassandra.allow_alter_rf_during_range_movement", "true");

System.setProperty("cassandra.minimum_replication_factor", "1");
}

Expand Down
2 changes: 2 additions & 0 deletions cassandra-analytics-integration-tests/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,8 @@ def configureIntegrationTestTask = { Test task, String majorMinor = null ->
// use testCassandra40/testCassandra41 for targeted runs. CI covers all versions via CASSANDRA_VERSION.
task.systemProperty "cassandra.sidecar.versions_to_test", "5.0"
}
// Dtest installs a SecurityManager; its default limit of 25 UDP sockets is too small for multi-node tests.
task.systemProperty "sun.net.maxDatagramSockets", "1024"
task.systemProperty "SKIP_STARTUP_VALIDATIONS", "true"
task.systemProperty "logback.configurationFile", "src/test/resources/logback-test.xml"
task.systemProperty "cassandra.integration.sidecar.test.enable_mtls", integrationEnableMtls
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,31 +28,40 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.IntStream;

import com.google.common.collect.Range;

import org.junit.jupiter.api.AfterAll;

import org.apache.cassandra.bridge.CassandraVersion;
import org.apache.cassandra.distributed.api.ConsistencyLevel;
import org.apache.cassandra.distributed.api.ICluster;
import org.apache.cassandra.distributed.api.IInstance;
import org.apache.cassandra.distributed.api.IInstanceConfig;
import org.apache.cassandra.distributed.api.Row;
import org.apache.cassandra.distributed.api.SimpleQueryResult;
import org.apache.cassandra.distributed.api.TokenSupplier;
import org.apache.cassandra.sidecar.common.server.JmxClient;
import org.apache.cassandra.sidecar.testing.QualifiedName;
import org.apache.cassandra.spark.bulkwriter.DecoratedKey;
import org.apache.cassandra.spark.bulkwriter.Tokenizer;
import org.apache.cassandra.spark.common.schema.ColumnType;
import org.apache.cassandra.spark.common.schema.ColumnTypes;
import org.apache.cassandra.testing.ClusterBuilderConfiguration;
import org.apache.cassandra.testing.TestTokenSupplier;
import org.apache.cassandra.testing.utils.WithProperties;
import scala.Tuple2;

import static org.apache.cassandra.testing.TestUtils.TEST_KEYSPACE;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assumptions.assumeTrue;
import static org.junit.jupiter.api.Assumptions.assumeFalse;

/**
* Base class for resiliency tests. Contains helper methods for data generation and validation
Expand All @@ -61,35 +70,91 @@ public abstract class ResiliencyTestBase extends SharedClusterSparkIntegrationTe
{
public static final String QUERY_ALL_ROWS = "SELECT * FROM %s";

private final WithProperties topologyChangeProperties = new WithProperties();

/**
* Skips the test class when the version under test declares none of the internals that the ByteBuddy hooks of a
* topology-change test intercept: {@code StorageService.bootstrap(Collection, long)},
* {@code StorageService.unbootstrap()} and {@code RangeRelocator.stream()}. CEP-21 Transactional Cluster Metadata
* removed all three in Cassandra 6.0, in favour of {@code org.apache.cassandra.tcm.sequences}. A hook that fails
* to install is silent, so the test would instead wait on a latch that never counts down.
*
* <p>CASSANALYTICS-112 tracks the port of these hooks to {@code BootstrapAndJoin.bootstrap},
* {@code UnbootstrapStreams.execute} and {@code Move.executeNext}, and the decision on the tests that move
* several nodes at once, which Transactional Cluster Metadata no longer permits.</p>
*
* <p>Call this from {@link #beforeClusterProvisioning()}, which runs before the cluster starts.</p>
* Legacy scenarios that need multiple nodes to remain in transition at once.
*/
protected void assumeTopologyChangeHooksSupported()
protected boolean requiresConcurrentTopologyChanges()
{
return false;
}

@Override
protected void beforeClusterProvisioning()
{
super.beforeClusterProvisioning();
if (requiresConcurrentTopologyChanges())
{
assumeFalse(usesTcm(), "Concurrent topology scenarios require Cassandra older than 6.0");
// Preserve the permissive configuration used by the legacy concurrent scenarios.
topologyChangeProperties.with("cassandra.consistent.rangemovement", "false",
"cassandra.consistent.simultaneousmoves.allow", "true",
"cassandra.allow_alter_rf_during_range_movement", "true");
}
}

@Override
@AfterAll
protected void tearDown() throws Exception
{
try
{
super.tearDown();
}
finally
{
topologyChangeProperties.close();
}
}

protected boolean usesTcm()
{
return CassandraVersion.fromVersion(testVersion.version())
.orElseThrow(() -> new IllegalStateException("Unsupported Cassandra version: " + testVersion))
.versionNumber() >= CassandraVersion.SIXZERO.versionNumber();
}

protected void prepareTopologyChange()
{
if (usesTcm())
{
// Keep a metadata quorum when a test stops or replaces a node.
cluster.get(1).nodetoolResult("cms", "reconfigure", "3").asserts().success();
}
}

protected static TokenSupplier disjointMultiDcTokens()
{
TokenSupplier tokens = TestTokenSupplier.evenlyDistributedTokens(6, 0, 2, 1);
// Match Sidecar CASSSIDECAR-277's swap(5, 10), whose indices are zero-based.
// Nodes 11 and 12 can then join or leave concurrently without overlapping range locks.
return node -> tokens.tokens(node == 6 ? 11 : node == 11 ? 6 : node);
}

protected static <T> T awaitTopologyChange(Future<T> task, boolean failureExpected)
{
String version = testVersion.version();
CassandraVersion underTest = CassandraVersion.fromVersion(version)
.orElseThrow(() -> new IllegalStateException(
"Unsupported Cassandra version for topology-change tests: " + version));
boolean supported = underTest.versionNumber() < CassandraVersion.SIXZERO.versionNumber();
if (!supported)
try
{
return task.get(2, TimeUnit.MINUTES);
}
catch (ExecutionException exception)
{
if (!failureExpected)
{
throw new AssertionError("Topology change failed", exception.getCause());
}
return null;
}
catch (InterruptedException exception)
{
Thread.currentThread().interrupt();
throw new AssertionError("Interrupted waiting for topology change", exception);
}
catch (TimeoutException exception)
{
// An aborted @BeforeAll produces no test event, so Gradle reports nothing
logger.warn("Skipping {}: a topology-change test intercepts Cassandra internals that CEP-21 removed "
+ "in 6.0, and the test version is {}", getClass().getSimpleName(), version);
throw new AssertionError("Topology change did not finish", exception);
}
assumeTrue(supported,
"A topology-change test intercepts Cassandra internals that CEP-21 removed in 6.0, "
+ "but the test version is " + version);
}

public Set<String> getDataForRange(Range<BigInteger> range, int rowCount)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -256,9 +256,13 @@ public void assertExpectedBulkWriteFailure(String writeCL, DataFrameWriter<Row>
cause = cause.getCause();
}

assertThat(cause).isNotNull()
.hasMessageFindingMatch("Failed to write (\\d+) ranges with " + writeCL +
" for job ([a-zA-Z0-9-]+) in phase .*");
if (cause == null)
{
throw new AssertionError("Expected a range write failure with " + writeCL, thrown);
}

assertThat(cause).hasMessageFindingMatch("Failed to write (\\d+) ranges with " + writeCL +
" for job ([a-zA-Z0-9-]+) in phase .*");
}

/**
Expand Down
Loading
Loading