[CASSANDRA-20736] Test fixes

This commit is contained in:
Sam Tunnicliffe 2025-06-16 19:18:56 +01:00
parent 360facb24c
commit 4af3e47b2b
11 changed files with 64 additions and 41 deletions

View File

@ -233,7 +233,13 @@ public class ClusterMetadata
public boolean isCMSMember(InetAddressAndPort endpoint)
{
return fullCMSMembers().contains(endpoint);
if (epoch.isAfter(Epoch.FIRST))
return fullCMSMembers().contains(endpoint);
// special case to handle initialization of the CMS for the first time
return epoch.isEqualOrBefore(Epoch.FIRST) && cmsDataPlacement.reads.byEndpoint()
.keySet()
.contains(endpoint);
}
public Set<InetAddressAndPort> fullCMSMembers()
@ -266,11 +272,6 @@ public class ClusterMetadata
return fullCMSReplicas;
}
public DataPlacement getCMSPlacement()
{
return cmsDataPlacement;
}
private DataPlacement calculateCMSPlacement(DataPlacements placements, CMSMembership cms)
{
if (epoch.isBefore(Epoch.FIRST) || schema.getKeyspaces().get(SchemaConstants.METADATA_KEYSPACE_NAME).isEmpty())

View File

@ -66,10 +66,7 @@ public class PaxosBackedProcessor extends AbstractLocalProcessor
@Override
protected boolean acceptCommit(ClusterMetadata metadata)
{
if (metadata.epoch.isAfter(Epoch.FIRST) && metadata.fullCMSMembers().contains(FBUtilities.getBroadcastAddressAndPort()))
return true;
return metadata.epoch.isEqualOrBefore(Epoch.FIRST)
&& metadata.getCMSPlacement().reads.byEndpoint().keySet().contains(FBUtilities.getBroadcastAddressAndPort());
return metadata.isCMSMember(FBUtilities.getBroadcastAddressAndPort());
}
@Override

View File

@ -32,14 +32,11 @@ import org.apache.cassandra.distributed.api.ConsistencyLevel;
import org.apache.cassandra.distributed.test.TestBaseImpl;
import org.apache.cassandra.io.util.FileOutputStreamPlus;
import org.apache.cassandra.locator.InetAddressAndPort;
import org.apache.cassandra.locator.MetaStrategy;
import org.apache.cassandra.locator.Replica;
import org.apache.cassandra.schema.ReplicationParams;
import org.apache.cassandra.tcm.CMSOperations;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.tcm.Epoch;
import org.apache.cassandra.tcm.membership.NodeId;
import org.apache.cassandra.tcm.membership.NodeVersion;
import org.apache.cassandra.tcm.ownership.DataPlacement;
import org.apache.cassandra.tcm.serialization.VerboseMetadataSerializer;
import static org.apache.cassandra.distributed.shared.ClusterUtils.start;
@ -108,15 +105,10 @@ public class BootWithMetadataTest extends TestBaseImpl
try
{
ClusterMetadata metadata = ClusterMetadata.current();
Replica oldCMS = MetaStrategy.replica(InetAddressAndPort.getByNameUnchecked("127.0.0.1"));
Replica newCMS = MetaStrategy.replica(InetAddressAndPort.getByNameUnchecked("127.0.0.2"));
NodeId oldCMS = metadata.directory.peerId(InetAddressAndPort.getByNameUnchecked("127.0.0.1"));
NodeId newCMS = metadata.directory.peerId(InetAddressAndPort.getByNameUnchecked("127.0.0.2"));
ClusterMetadata.Transformer transformer = metadata.transformer();
DataPlacement.Builder builder = metadata.placements.get(ReplicationParams.meta(metadata)).unbuild()
.withoutReadReplica(metadata.nextEpoch(), oldCMS)
.withoutWriteReplica(metadata.nextEpoch(), oldCMS)
.withWriteReplica(metadata.nextEpoch(), newCMS)
.withReadReplica(metadata.nextEpoch(), newCMS);
transformer = transformer.with(metadata.placements.unbuild().with(ReplicationParams.meta(metadata), builder.build()).build());
transformer.leaveCMS(oldCMS).startJoiningCMS(newCMS).finishJoiningCMS(newCMS);
ClusterMetadata toDump = transformer.build().metadata.forceEpoch(Epoch.create(1000));
Path p = Files.createTempFile("clustermetadata", "dump");
try (FileOutputStreamPlus out = new FileOutputStreamPlus(p))

View File

@ -66,7 +66,8 @@ public class ClusterMetadataDumpTest extends TestBaseImpl
epochsSeen++;
}
assertEquals(3, unsafeJoinSeen);
assertEquals(3, registerSeen);
// Only 2 REGISTER transforms are expected as the first CMS node is registered implicitly by INITIALIZE_CMS
assertEquals(2, registerSeen);
assertTrue(epochsSeen > 15);
res = cluster.get(1).nodetoolResult("cms", "dumplog", "--start", "10", "--end", "15");

View File

@ -41,7 +41,6 @@ import org.apache.cassandra.distributed.shared.NetworkTopology;
import org.apache.cassandra.locator.MetaStrategy;
import org.apache.cassandra.schema.DistributedMetadataLogKeyspace;
import org.apache.cassandra.schema.ReplicationParams;
import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.service.paxos.Ballot;
import org.apache.cassandra.service.paxos.PaxosRepairHistory;
import org.apache.cassandra.tcm.ClusterMetadata;
@ -57,6 +56,7 @@ import static org.apache.cassandra.distributed.shared.ClusterUtils.awaitRingJoin
import static org.apache.cassandra.distributed.shared.ClusterUtils.replaceHostAndStart;
import static org.apache.cassandra.distributed.shared.NetworkTopology.dcAndRack;
import static org.apache.cassandra.distributed.shared.NetworkTopology.networkTopology;
import static org.apache.cassandra.schema.SchemaConstants.METADATA_KEYSPACE_NAME;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.psjava.util.AssertStatus.assertTrue;
@ -85,10 +85,10 @@ public class ReconfigureCMSTest extends FuzzTestBase
ClusterMetadata metadata = ClusterMetadata.current();
assertEquals(5, metadata.fullCMSMembers().size());
assertEquals(ReplicationParams.simpleMeta(5, metadata.directory.knownDatacenters()),
metadata.placements().keys().stream().filter(ReplicationParams::isMeta).findFirst().get());
metadata.schema.getKeyspaceMetadata(METADATA_KEYSPACE_NAME).params.replication);
});
cluster.stream().forEach(i -> {
Assert.assertTrue(i.executeInternal(String.format("SELECT * FROM %s.%s", SchemaConstants.METADATA_KEYSPACE_NAME, DistributedMetadataLogKeyspace.TABLE_NAME)).length > 0);
Assert.assertTrue(i.executeInternal(String.format("SELECT * FROM %s.%s", METADATA_KEYSPACE_NAME, DistributedMetadataLogKeyspace.TABLE_NAME)).length > 0);
});
cluster.get(nodeSelector.get()).nodetoolResult("cms", "reconfigure", "1").asserts().success();
@ -96,7 +96,7 @@ public class ReconfigureCMSTest extends FuzzTestBase
ClusterMetadata metadata = ClusterMetadata.current();
assertEquals(1, metadata.fullCMSMembers().size());
assertEquals(ReplicationParams.simpleMeta(1, metadata.directory.knownDatacenters()),
metadata.placements().keys().stream().filter(ReplicationParams::isMeta).findFirst().get());
metadata.schema.getKeyspaceMetadata(METADATA_KEYSPACE_NAME).params.replication);
});
}
}
@ -326,11 +326,11 @@ public class ReconfigureCMSTest extends FuzzTestBase
Object[][] rows = instance.executeInternal("select points from system.paxos_repair_history " +
"where keyspace_name = ? " +
"and table_name = ?",
SchemaConstants.METADATA_KEYSPACE_NAME,
METADATA_KEYSPACE_NAME,
DistributedMetadataLogKeyspace.TABLE_NAME);
if (rows.length == 0)
return PaxosRepairHistory.empty(SchemaConstants.METADATA_KEYSPACE_NAME, DistributedMetadataLogKeyspace.TABLE_NAME);
return PaxosRepairHistory.empty(METADATA_KEYSPACE_NAME, DistributedMetadataLogKeyspace.TABLE_NAME);
assertEquals(1, rows.length);
//noinspection unchecked
List<ByteBuffer> points = (List<ByteBuffer>)rows[0][0];

View File

@ -18,6 +18,7 @@
package org.apache.cassandra.auth;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
import java.util.Set;
@ -34,6 +35,8 @@ import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.cql3.CQLTester;
import org.apache.cassandra.db.ConsistencyLevel;
import org.apache.cassandra.db.SystemKeyspace;
import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper;
import org.apache.cassandra.schema.ReplicationParams;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.schema.TableMetadata;
@ -59,6 +62,11 @@ public class GrantAndRevokeTest extends CQLTester
ServerTestUtils.daemonInitialization();
DatabaseDescriptor.setPermissionsValidity(0);
DatabaseDescriptor.setRolesValidity(0);
// for the tables in the distributed metadata keyspace to be queryable, there needs to be a valid placement
// which is derived from the CMS membership. Most unit tests don't actually need to query those dist tables
// which is why this isn't done as a matter of routine
ClusterMetadataTestHelper.reconfigureCms(ReplicationParams.ntsMeta(Collections.singletonMap(DatabaseDescriptor.getLocalDataCenter(), 1)));
ServerTestUtils.markCMS();
requireAuthentication();
requireNetwork();
CassandraDaemon.getInstanceForTesting().setupVirtualKeyspaces();

View File

@ -41,6 +41,7 @@ import org.apache.cassandra.io.util.FileOutputStreamPlus;
import org.apache.cassandra.tcm.membership.Directory;
import org.apache.cassandra.tcm.membership.Location;
import org.apache.cassandra.tcm.membership.MembershipUtils;
import org.apache.cassandra.tcm.membership.NodeAddresses;
import org.apache.cassandra.tcm.membership.NodeId;
import org.apache.cassandra.tcm.membership.NodeVersion;
import org.apache.cassandra.tcm.ownership.DataPlacements;
@ -110,7 +111,10 @@ public class BootWithMetadataTest
Directory directory = first.directory;
int nodeCount = 10;
int tokensPerNode = 5;
for (int i = 0; i < nodeCount; i++)
// Ensure that the "local" node is registered as it being member of the CMS is a precondition of booting from
// a ClusterMetadata.
directory = directory.with(NodeAddresses.current(), new Location("DC1", "RACK1"));
for (int i = 0; i < nodeCount-1; i++)
directory = directory.with(nodeAddresses(random), new Location("DC1", "RACK1"));
t = t.with(directory);

View File

@ -41,8 +41,12 @@ import org.apache.cassandra.tcm.log.LocalLog;
import org.apache.cassandra.tcm.log.LogState;
import org.apache.cassandra.tcm.log.LogStorage;
import org.apache.cassandra.tcm.log.SystemKeyspaceStorage;
import org.apache.cassandra.tcm.membership.NodeAddresses;
import org.apache.cassandra.tcm.membership.NodeId;
import org.apache.cassandra.tcm.membership.NodeVersion;
import org.apache.cassandra.tcm.ownership.UniformRangePlacement;
import org.apache.cassandra.tcm.transformations.CustomTransformation;
import org.apache.cassandra.tcm.transformations.ForceSnapshot;
import org.apache.cassandra.utils.FBUtilities;
import static org.junit.Assert.assertEquals;
@ -80,6 +84,16 @@ public class GetLogStateTest
ClusterMetadataService.setInstance(cms);
log.readyUnchecked();
log.unsafeBootstrapForTesting(FBUtilities.getBroadcastAddressAndPort());
// register first node & make it CMS as this affects how the meta strategy placements are calculated
ClusterMetadata withRegistered = ClusterMetadata.current()
.transformer()
.register(NodeAddresses.current(),
DatabaseDescriptor.getInitialLocationProvider().initialLocation(),
NodeVersion.CURRENT)
.build().metadata;
NodeId id = withRegistered.directory.peerId(NodeAddresses.current().broadcastAddress);
ClusterMetadata withCMS = withRegistered.transformer().startJoiningCMS(id).finishJoiningCMS(id).build().metadata;
cms.commit(new ForceSnapshot(withCMS));
}
@Test

View File

@ -62,8 +62,9 @@ public class DistributedLogStateTest extends LogStateTestBase
{
return new LogStateSUT()
{
// start test entries at FIRST + 1 as the pre-init transform is automatically inserted with Epoch.FIRST
// we start test entries at FIRST, but in a real log the PRE_INITIALIZE_CMS transform is automatically
// inserted with Epoch.FIRST, followed by INITIALIZE_CMS so the next entry to be committed would be at
// epoch 3
Epoch currentEpoch = Epoch.FIRST;
Epoch nextEpoch;
boolean applied;

View File

@ -105,11 +105,13 @@ public class InProgressSequenceCancellationTest
LockedRanges locked = lockedRanges(placements, random);
// state of metadata before starting the sequence
// note: don't allow epoch to be Epoch.FIRST as this is a
// special case for calculating meta strategy placements.
ClusterMetadata before = metadata(directory).transformer()
.with(placements)
.withNodeState(nodeId, NodeState.REGISTERED)
.with(locked)
.build().metadata;
.build().metadata.forceEpoch(epoch(random));
// Placements after PREPARE_JOIN
DataPlacements afterPrepare = placements(ranges(random), replication, random);
@ -179,11 +181,13 @@ public class InProgressSequenceCancellationTest
// Ranges locked by other operations
LockedRanges locked = lockedRanges(placements, random);
// state of metadata before starting the sequence
// note: don't allow epoch to be Epoch.FIRST as this is a
// special case for calculating meta strategy placements.
ClusterMetadata before = metadata(directory).transformer()
.with(placements)
.withNodeState(nodeId, NodeState.JOINED)
.with(locked)
.build().metadata;
.build().metadata.forceEpoch(epoch(random));
// PREPARE_LEAVE does not modify placements, so first transformation is START_LEAVE

View File

@ -84,10 +84,11 @@ public class PrepareLeaveTest
public void testCheckRF_Simple() throws Throwable
{
Keyspaces kss = Keyspaces.of(DistributedMetadataLogKeyspace.initialMetadata(Sets.newHashSet(hostDc.values())), KSM);
// should be accepted (2 nodes in dc1 where we remove the host):
ClusterMetadata metadata = prepMetadata(kss, 2, 2);
assertTrue(executeLeave(metadata));
// should be rejected:
metadata = prepMetadata(kss, 1, 2);
// should be rejected because only 1 node in dc2:
metadata = prepMetadata(kss, 2, 1);
assertFalse(executeLeave(metadata));
}
@ -97,17 +98,17 @@ public class PrepareLeaveTest
Keyspaces kss = Keyspaces.of(DistributedMetadataLogKeyspace.initialMetadata(Sets.newHashSet(hostDc.values())), KSM_NTS);
ClusterMetadata metadata = prepMetadata(kss, 4, 4);
assertTrue(executeLeave(metadata));
// should be accepted (4 nodes in dc1 where we remove the host):
metadata = prepMetadata(kss, 4, 2);
// should be accepted (4 nodes in dc2 where we remove the host):
metadata = prepMetadata(kss, 2, 4);
assertTrue(executeLeave(metadata));
// should be rejected
metadata = prepMetadata(kss, 3, 4);
// should be rejected because there are already only 3 replicas in dc2
metadata = prepMetadata(kss, 4, 3);
assertFalse(executeLeave(metadata));
}
private boolean executeLeave(ClusterMetadata metadata) throws Throwable
{
PrepareLeave prepareLeave = new PrepareLeave(metadata.directory.peerId(InetAddressAndPort.getByName("127.0.0.1")),
PrepareLeave prepareLeave = new PrepareLeave(metadata.directory.peerId(InetAddressAndPort.getByName("127.0.0.11")),
false,
dummyPlacementProvider,
LeaveStreams.Kind.UNBOOTSTRAP);