diff --git a/src/java/org/apache/cassandra/tcm/ClusterMetadata.java b/src/java/org/apache/cassandra/tcm/ClusterMetadata.java index 23c06a1227..423d2c82af 100644 --- a/src/java/org/apache/cassandra/tcm/ClusterMetadata.java +++ b/src/java/org/apache/cassandra/tcm/ClusterMetadata.java @@ -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 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()) diff --git a/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java b/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java index 80b11d66f8..be974c1247 100644 --- a/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java +++ b/src/java/org/apache/cassandra/tcm/PaxosBackedProcessor.java @@ -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 diff --git a/test/distributed/org/apache/cassandra/distributed/test/log/BootWithMetadataTest.java b/test/distributed/org/apache/cassandra/distributed/test/log/BootWithMetadataTest.java index bcc02b89a6..a76d0d8d42 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/log/BootWithMetadataTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/log/BootWithMetadataTest.java @@ -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)) diff --git a/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java b/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java index 7731bde594..c0a4466e78 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/log/ClusterMetadataDumpTest.java @@ -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"); diff --git a/test/distributed/org/apache/cassandra/distributed/test/log/ReconfigureCMSTest.java b/test/distributed/org/apache/cassandra/distributed/test/log/ReconfigureCMSTest.java index 64663280aa..e2951a2f84 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/log/ReconfigureCMSTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/log/ReconfigureCMSTest.java @@ -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 points = (List)rows[0][0]; diff --git a/test/unit/org/apache/cassandra/auth/GrantAndRevokeTest.java b/test/unit/org/apache/cassandra/auth/GrantAndRevokeTest.java index f0060dae93..58c3038229 100644 --- a/test/unit/org/apache/cassandra/auth/GrantAndRevokeTest.java +++ b/test/unit/org/apache/cassandra/auth/GrantAndRevokeTest.java @@ -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(); diff --git a/test/unit/org/apache/cassandra/tcm/BootWithMetadataTest.java b/test/unit/org/apache/cassandra/tcm/BootWithMetadataTest.java index 476a454597..fe9c6f0264 100644 --- a/test/unit/org/apache/cassandra/tcm/BootWithMetadataTest.java +++ b/test/unit/org/apache/cassandra/tcm/BootWithMetadataTest.java @@ -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); diff --git a/test/unit/org/apache/cassandra/tcm/GetLogStateTest.java b/test/unit/org/apache/cassandra/tcm/GetLogStateTest.java index 167aa3992b..965b3f4918 100644 --- a/test/unit/org/apache/cassandra/tcm/GetLogStateTest.java +++ b/test/unit/org/apache/cassandra/tcm/GetLogStateTest.java @@ -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 diff --git a/test/unit/org/apache/cassandra/tcm/log/DistributedLogStateTest.java b/test/unit/org/apache/cassandra/tcm/log/DistributedLogStateTest.java index a4c884406c..fc8c0559e1 100644 --- a/test/unit/org/apache/cassandra/tcm/log/DistributedLogStateTest.java +++ b/test/unit/org/apache/cassandra/tcm/log/DistributedLogStateTest.java @@ -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; diff --git a/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java b/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java index 94eb01fc1b..50946d7ac4 100644 --- a/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java +++ b/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java @@ -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 diff --git a/test/unit/org/apache/cassandra/tcm/transformations/PrepareLeaveTest.java b/test/unit/org/apache/cassandra/tcm/transformations/PrepareLeaveTest.java index f5044a00fd..bbf58eda54 100644 --- a/test/unit/org/apache/cassandra/tcm/transformations/PrepareLeaveTest.java +++ b/test/unit/org/apache/cassandra/tcm/transformations/PrepareLeaveTest.java @@ -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);