result = fetcher.asyncFetchLog(REMOTE, Epoch.FIRST);
+ // The continuation has run by the time asyncFetchLog returns (inline on the caller for
+ // the buggy code, on the responder thread otherwise). Wait for the future to settle -
+ // it completes exceptionally here since the fetched epoch never reaches FIRST, which is
+ // irrelevant: we only care which thread ran the continuation.
+ result.awaitUninterruptibly(30, TimeUnit.SECONDS);
+
+ Thread ran = continuationThread.get();
+ assertNotNull("The append/wait continuation never ran", ran);
+ assertNotSame("CASSANDRA-21384: the peer-log append/wait continuation must not run on the " +
+ "caller (e.g. GossipStage) thread, otherwise it can deadlock with GlobalLogFollower",
+ callerThread, ran);
+ assertSame("The continuation should run on the messaging response thread",
+ responderThread.get(), ran);
+ }
+ finally
+ {
+ MessagingService.instance().outboundSink.remove(interceptor);
+ responder.shutdownNow();
+ log.close();
+ ClusterMetadataService.unsetInstance();
+ }
+ }
+
+ /**
+ * Regression test for CASSANDRA-21384 that reproduces the *actual* circular wait from the
+ * ticket end-to-end (not just the thread on which the continuation runs), and verifies it with
+ * a thread dump.
+ *
+ * It runs the peer-log fetch on the real, single-threaded {@link Stage#GOSSIP} and registers a
+ * post-commit {@link ChangeListener} that - like the production listeners - calls
+ * {@link Gossiper#runInGossipStageBlocking} while an epoch is enacted. On the pre-fix code the
+ * fetch's append/wait continuation runs inline on the GossipStage thread and blocks it, so:
+ *
+ *
+ * GossipStage : blocked in PeerLogFetcher -> LocalLog.waitForHighestConsecutive()
+ * (waiting for the log to advance)
+ * GlobalLogFollower : enacting the epoch -> ChangeListener -> Gossiper.runInGossipStageBlocking()
+ * (waiting for the single GossipStage thread, which is stuck above)
+ *
+ *
+ * The test asserts the fetch returns within a timeout (i.e. no deadlock). If it does deadlock
+ * it captures a full thread dump and asserts the two stacks are exactly that cycle before
+ * failing. On the fixed code the continuation runs on the messaging thread, GossipStage returns
+ * from the fetch immediately, and the follower's task is serviced - the test passes.
+ */
+ @Test
+ public void doesNotDeadlockBetweenGossipStageAndLogFollower() throws Exception
+ {
+ LocalLog log = LocalLog.logSpec()
+ .async()
+ .withInitialState(new ClusterMetadata(Murmur3Partitioner.instance))
+ .createLog();
+ ClusterMetadataService cms = new ClusterMetadataService(new UniformRangePlacement(),
+ MetadataSnapshots.NO_OP,
+ log,
+ new AtomicLongBackedProcessor(log),
+ Commit.Replicator.NO_OP,
+ false);
+ ClusterMetadataService.setInstance(cms);
+ log.readyUnchecked();
+
+ // Models the production post-commit listeners (e.g. LegacyStateListener) which block on the
+ // gossip stage while the follower enacts an epoch. Only armed for the fetch under test.
+ AtomicBoolean armed = new AtomicBoolean(false);
+ log.addListener(new ChangeListener()
+ {
+ @Override
+ public void notifyPostCommit(ClusterMetadata prev, ClusterMetadata next, boolean fromSnapshot)
+ {
+ if (armed.get())
+ Gossiper.runInGossipStageBlocking(() -> {});
+ }
+ });
+
+ // Releases the caller from the intercepted send once the request promise has been completed
+ // (buggy path) or the continuation has begun appending on the responder thread (fixed path).
+ // The latter frees the GossipStage thread before the follower needs it, so the fixed code
+ // does not deadlock.
+ CountDownLatch releaseSend = new CountDownLatch(1);
+ log.addFilter(entry -> {
+ releaseSend.countDown();
+ return false; // keep the entry so the follower enacts it and fires the listener
+ });
+
+ ExecutorService responder = Executors.newSingleThreadExecutor(r -> new Thread(r, "test-fetch-responder"));
+
+ // Complete the request promise from the responder thread (mimicking the messaging response
+ // callback), holding the caller inside the send until that has happened - so the buggy
+ // code's map() sees an already-completed promise and runs the continuation inline on the
+ // GossipStage thread. The message is always dropped so nothing hits the network.
+ OutboundSink.Filter interceptor = (message, to, type) -> {
+ if (message.verb() != Verb.TCM_FETCH_PEER_LOG_REQ)
+ return true;
+
+ long id = message.id();
+ Message> request = message;
+ responder.execute(() -> {
+ // A snapshot at FIRST that the follower will enact, advancing the epoch and firing
+ // the post-commit listener above.
+ LogState logState = LogState.make(new ClusterMetadata(Murmur3Partitioner.instance).forceEpoch(Epoch.FIRST));
+ Message response = request.responseWith(logState);
+ MessagingService.instance().callbacks.removeAndRespond(id, to, response);
+ releaseSend.countDown();
+ });
+ Uninterruptibles.awaitUninterruptibly(releaseSend, 30, TimeUnit.SECONDS);
+ return false;
+ };
+ MessagingService.instance().outboundSink.add(interceptor);
+
+ AtomicReference> fetchFuture = new AtomicReference<>();
+ CountDownLatch fetchReturned = new CountDownLatch(1);
+ try
+ {
+ PeerLogFetcher fetcher = new PeerLogFetcher(log);
+ armed.set(true);
+ // Run the fetch on the real, single-threaded GossipStage, exactly as the gossip path does.
+ Stage.GOSSIP.execute(() -> {
+ try
+ {
+ fetchFuture.set(fetcher.asyncFetchLog(REMOTE, Epoch.FIRST));
+ }
+ finally
+ {
+ fetchReturned.countDown();
+ }
+ });
+
+ if (!fetchReturned.await(20, TimeUnit.SECONDS))
+ {
+ // GossipStage never returned from the fetch: the deadlock. Prove it is exactly the
+ // GlobalLogFollower <-> GossipStage cycle, then fail with the dump attached.
+ String dump = fullThreadDump();
+ StackTraceElement[] gossip = stackOf("GossipStage");
+ StackTraceElement[] follower = stackOf("GlobalLogFollower");
+
+ // break the deadlock so we don't leak the shared GossipStage thread
+ interruptThreads("GossipStage");
+
+ assertTrue("Expected GossipStage to be blocked advancing the log inside the peer-log " +
+ "fetch continuation, but was:\n" + dump,
+ stackContains(gossip, "PeerLogFetcher") && stackContains(gossip, "LocalLog"));
+ assertTrue("Expected GlobalLogFollower to be blocked handing off to the gossip stage, " +
+ "but was:\n" + dump,
+ stackContains(follower, "runInGossipStageBlocking"));
+ fail("CASSANDRA-21384: reproduced deadlock between GlobalLogFollower and GossipStage:\n" + dump);
+ }
+
+ // No deadlock: the fetch returned on the gossip stage. Let it settle to be sure the
+ // whole chain (responder -> follower -> gossip stage) completed.
+ Future f = fetchFuture.get();
+ assertNotNull("peer-log fetch future was not captured", f);
+ f.awaitUninterruptibly(20, TimeUnit.SECONDS);
+ }
+ finally
+ {
+ armed.set(false);
+ MessagingService.instance().outboundSink.remove(interceptor);
+ responder.shutdownNow();
+ log.close();
+ ClusterMetadataService.unsetInstance();
+ }
+ }
+
+ private static StackTraceElement[] stackOf(String namePrefix)
+ {
+ for (Map.Entry e : Thread.getAllStackTraces().entrySet())
+ if (e.getKey().getName().startsWith(namePrefix))
+ return e.getValue();
+ return null;
+ }
+
+ private static boolean stackContains(StackTraceElement[] stack, String needle)
+ {
+ if (stack == null)
+ return false;
+ for (StackTraceElement frame : stack)
+ if (frame.toString().contains(needle))
+ return true;
+ return false;
+ }
+
+ private static void interruptThreads(String namePrefix)
+ {
+ for (Thread t : Thread.getAllStackTraces().keySet())
+ if (t.getName().startsWith(namePrefix))
+ t.interrupt();
+ }
+
+ private static String fullThreadDump()
+ {
+ StringBuilder sb = new StringBuilder();
+ for (Map.Entry e : Thread.getAllStackTraces().entrySet())
+ {
+ Thread t = e.getKey();
+ sb.append('"').append(t.getName()).append("\" ").append(t.getState()).append('\n');
+ for (StackTraceElement frame : e.getValue())
+ sb.append("\tat ").append(frame).append('\n');
+ sb.append('\n');
+ }
+ return sb.toString();
+ }
+}
diff --git a/test/unit/org/apache/cassandra/tcm/UnregisterTest.java b/test/unit/org/apache/cassandra/tcm/UnregisterTest.java
index 62a6b160a2..3861eefd6b 100644
--- a/test/unit/org/apache/cassandra/tcm/UnregisterTest.java
+++ b/test/unit/org/apache/cassandra/tcm/UnregisterTest.java
@@ -110,7 +110,7 @@ public class UnregisterTest
assertFalse(metadata.directory.allJoinedEndpoints().contains(ep));
assertFalse(metadata.directory.allDatacenterRacks().containsKey("dc2"));
assertFalse(metadata.directory.knownDatacenters().contains("dc2"));
- metadata.placements.asMap().forEach((params, placement) -> {
+ metadata.placements().forEach((params, placement) -> {
assertFalse(Streams.concat(placement.writes.endpoints.stream(), placement.reads.endpoints.stream()).anyMatch((fr) -> fr.endpoints().contains(ep)));
});
}
diff --git a/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java b/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java
index cac53d553f..17d515e1ff 100644
--- a/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java
+++ b/test/unit/org/apache/cassandra/tcm/compatibility/GossipHelperTest.java
@@ -104,7 +104,7 @@ public class GossipHelperTest
assertEquals(internal, metadata.directory.addresses.get(nodeId).localAddress);
assertEquals(nativeAddress, metadata.directory.addresses.get(nodeId).nativeAddress);
- DataPlacements dp = metadata.placements;
+ DataPlacements dp = metadata.placements();
assertEquals(1, dp.get(KSM.params.replication).reads.forToken(token).get().size());
assertTrue(dp.get(KSM.params.replication).reads.forToken(token).get().contains(endpoint));
assertEquals(1, dp.get(KSM.params.replication).writes.forToken(token).get().size());
@@ -196,8 +196,8 @@ public class GossipHelperTest
assertEquals(entry.getValue(), metadata.tokenMap.tokens(nodeId).iterator().next());
}
- ReplicaGroups reads = metadata.placements.get(KSM_NTS.params.replication).reads;
- ReplicaGroups writes = metadata.placements.get(KSM_NTS.params.replication).writes;
+ ReplicaGroups reads = metadata.placement(KSM_NTS.params.replication).reads;
+ ReplicaGroups writes = metadata.placement(KSM_NTS.params.replication).writes;
assertEquals(reads, writes);
// tokens are
// dc1: 1: 1000, 3: 3000, 5: 5000, 6: 7000, 7: 9000
diff --git a/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java b/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java
index 46adc7c67a..05bc00cadb 100644
--- a/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java
+++ b/test/unit/org/apache/cassandra/tcm/listeners/MetadataSnapshotListenerTest.java
@@ -109,7 +109,7 @@ public class MetadataSnapshotListenerTest
listener.notify(entry, result);
ClusterMetadata snapshot = snapshots.getSnapshot(nextEpoch);
assertEquals(nextEpoch, snapshot.epoch);
- assertEquals(toSnapshot.placements, snapshot.placements);
+ assertEquals(toSnapshot.placements(), snapshot.placements());
}
private MetadataSnapshots init()
diff --git a/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java b/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java
new file mode 100644
index 0000000000..70705ff4f9
--- /dev/null
+++ b/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java
@@ -0,0 +1,179 @@
+/*
+ * 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.tcm.listeners;
+
+import java.util.Random;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.junit.BeforeClass;
+import org.junit.Test;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import org.apache.cassandra.config.CassandraRelevantProperties;
+import org.apache.cassandra.dht.Murmur3Partitioner;
+import org.apache.cassandra.distributed.test.log.ClusterMetadataTestHelper;
+import org.apache.cassandra.locator.Replica;
+import org.apache.cassandra.schema.DistributedSchema;
+import org.apache.cassandra.schema.KeyspaceMetadata;
+import org.apache.cassandra.schema.KeyspaceParams;
+import org.apache.cassandra.schema.Keyspaces;
+import org.apache.cassandra.tcm.CMSMembership;
+import org.apache.cassandra.tcm.ClusterMetadata;
+import org.apache.cassandra.tcm.Epoch;
+import org.apache.cassandra.tcm.membership.MembershipUtils;
+import org.apache.cassandra.tcm.membership.NodeId;
+import org.apache.cassandra.tcm.ownership.DataPlacement;
+import org.apache.cassandra.tcm.ownership.DataPlacements;
+import org.apache.cassandra.tcm.ownership.OwnershipUtils;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotEquals;
+import static org.psjava.util.AssertStatus.assertTrue;
+
+public class PlacementsChangeListenerTest
+{
+ private static final Logger logger = LoggerFactory.getLogger(PlacementsChangeListenerTest.class);
+ static Random random;
+ static Epoch e;
+ static NodeId node1;
+ static KeyspaceMetadata ks1;
+
+ @BeforeClass
+ public static void setupClass()
+ {
+ long seed = System.nanoTime();
+ logger.info("Seed: {}", seed);
+ random = new Random(seed);
+ e = Epoch.create(10);
+ node1 = new NodeId(1);
+ ks1 = KeyspaceMetadata.create("ks1", KeyspaceParams.simple(1));
+ CassandraRelevantProperties.TCM_SORT_REPLICA_GROUPS.setBoolean(false);
+ }
+
+ @Test
+ public void testPlacementChange()
+ {
+ DataPlacements before = OwnershipUtils.randomPlacements(random).withLastModified(e);
+ DataPlacements.Builder builder = before.unbuild();
+ before.forEach((params, placement) -> {
+ Replica remove = placement.writes.byEndpoint().flattenValues().iterator().next();
+ Replica add = Replica.fullReplica(MembershipUtils.endpoint(99), remove.range());
+ DataPlacement newPlacement = placement.unbuild()
+ .withoutWriteReplica(e.nextEpoch(), remove)
+ .withWriteReplica(e.nextEpoch(), add).build();
+ builder.with(params, newPlacement);
+ });
+ DataPlacements after = builder.build().withLastModified(e.nextEpoch());
+
+ // only placements are different
+ ClusterMetadata prev = metadata(e, before, Keyspaces.of(ks1), node1);
+ ClusterMetadata next = metadata(e.nextEpoch(), after, Keyspaces.of(ks1), node1);
+ assertNotEquals(prev.placements().lastModified(), next.placements().lastModified());
+ assertFalse(prev.placements().equivalentTo(next.placements()));
+ assertEquals(prev.schema.getKeyspaces().size(), next.schema.getKeyspaces().size());
+ assertEquals(prev.schema.getKeyspaceMetadata("ks1").params, next.schema.getKeyspaceMetadata("ks1").params);
+ assertEquals(prev.cmsMembership, next.cmsMembership);
+
+ assertOnChangeEvent(prev, next);
+ }
+
+ @Test
+ public void testKeyspaceCountChange()
+ {
+ DataPlacements placements = OwnershipUtils.randomPlacements(random).withLastModified(e);
+ KeyspaceMetadata ks2 = KeyspaceMetadata.create("ks2", KeyspaceParams.simple(1));
+
+ // only keyspace counts are different
+ ClusterMetadata prev = metadata(e, placements, Keyspaces.of(ks1), node1);
+ ClusterMetadata next = metadata(e, placements, Keyspaces.of(ks1, ks2), node1);
+ assertEquals(prev.placements().lastModified(), next.placements().lastModified());
+ assertTrue(prev.placements().equivalentTo(next.placements()));
+ assertNotEquals(prev.schema.getKeyspaces().size(), next.schema.getKeyspaces().size());
+ assertEquals(prev.schema.getKeyspaceMetadata("ks1").params, next.schema.getKeyspaceMetadata("ks1").params);
+ assertEquals(prev.cmsMembership, next.cmsMembership);
+
+ assertOnChangeEvent(prev, next);
+ }
+
+ @Test
+ public void testKeyspaceParamsChange()
+ {
+ DataPlacements placements = OwnershipUtils.randomPlacements(random).withLastModified(e);
+ KeyspaceMetadata ks1a = KeyspaceMetadata.create("ks1", KeyspaceParams.simple(2));
+
+ // only keyspace params are different
+ ClusterMetadata prev = metadata(e, placements, Keyspaces.of(ks1), node1);
+ ClusterMetadata next = metadata(e, placements, Keyspaces.of(ks1a), node1);
+ assertEquals(prev.placements().lastModified(), next.placements().lastModified());
+ assertTrue(prev.placements().equivalentTo(next.placements()));
+ assertEquals(prev.schema.getKeyspaces().size(), next.schema.getKeyspaces().size());
+ assertNotEquals(prev.schema.getKeyspaceMetadata("ks1").params, next.schema.getKeyspaceMetadata("ks1").params);
+ assertEquals(prev.cmsMembership, next.cmsMembership);
+
+ assertOnChangeEvent(prev, next);
+ }
+
+ @Test
+ public void testCMSMembershipChange()
+ {
+ DataPlacements placements = OwnershipUtils.randomPlacements(random).withLastModified(e);
+ NodeId node2 = new NodeId(2);
+
+ // only cms memberships are different
+ ClusterMetadata prev = metadata(e, placements, Keyspaces.of(ks1), node1);
+ ClusterMetadata next = metadata(e, placements, Keyspaces.of(ks1), node1, node2);
+ assertEquals(prev.placements().lastModified(), next.placements().lastModified());
+ assertTrue(prev.placements().equivalentTo(next.placements()));
+ assertEquals(prev.schema.getKeyspaces().size(), next.schema.getKeyspaces().size());
+ assertEquals(prev.schema.getKeyspaceMetadata("ks1").params, next.schema.getKeyspaceMetadata("ks1").params);
+ assertNotEquals(prev.cmsMembership, next.cmsMembership);
+
+ assertOnChangeEvent(prev, next);
+ }
+
+ private static ClusterMetadata metadata(Epoch epoch,
+ DataPlacements placements,
+ Keyspaces keyspaces,
+ NodeId...cmsNode)
+ {
+ CMSMembership cms = CMSMembership.EMPTY;
+ for (NodeId n : cmsNode)
+ cms = cms.startJoining(n).finishJoining(n);
+
+ ClusterMetadata.Transformer t = ClusterMetadataTestHelper.minimalForTesting(epoch,
+ Murmur3Partitioner.instance,
+ new DistributedSchema(keyspaces, epoch),
+ cms)
+ .forceEpoch(epoch)
+ .transformer()
+ .with(placements);
+
+ return t.build().metadata;
+ }
+
+ private static void assertOnChangeEvent(ClusterMetadata prev, ClusterMetadata next)
+ {
+ AtomicInteger cnt = new AtomicInteger(0);
+ PlacementsChangeListener listener = new PlacementsChangeListener(cnt::incrementAndGet);
+ listener.notifyPostCommit(prev, next, false);
+ assertEquals(1, cnt.get());
+ }
+}
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/membership/MembershipUtils.java b/test/unit/org/apache/cassandra/tcm/membership/MembershipUtils.java
index 91e1f1e1c3..3820c8f8d4 100644
--- a/test/unit/org/apache/cassandra/tcm/membership/MembershipUtils.java
+++ b/test/unit/org/apache/cassandra/tcm/membership/MembershipUtils.java
@@ -19,8 +19,8 @@
package org.apache.cassandra.tcm.membership;
import java.net.UnknownHostException;
+import java.util.List;
import java.util.Random;
-import java.util.Set;
import java.util.stream.Collectors;
import org.apache.cassandra.locator.InetAddressAndPort;
@@ -38,13 +38,13 @@ public class MembershipUtils
return endpoint(random.nextInt(254) + 1);
}
- public static Set uniqueEndpoints(Random random, int count)
+ public static List uniqueEndpoints(Random random, int count)
{
return random.ints(1, 255)
.distinct()
.limit(count)
.mapToObj(MembershipUtils::endpoint)
- .collect(Collectors.toSet());
+ .collect(Collectors.toList());
}
public static InetAddressAndPort endpoint(int i)
diff --git a/test/unit/org/apache/cassandra/tcm/ownership/LocalRangesAllSettledTest.java b/test/unit/org/apache/cassandra/tcm/ownership/LocalRangesAllSettledTest.java
index f9e8104ca9..07a3dc4b03 100644
--- a/test/unit/org/apache/cassandra/tcm/ownership/LocalRangesAllSettledTest.java
+++ b/test/unit/org/apache/cassandra/tcm/ownership/LocalRangesAllSettledTest.java
@@ -97,7 +97,7 @@ public class LocalRangesAllSettledTest
AllLocalRanges proposed = snapshotAllLocalRanges(LocalRangeStatus.SETTLED, INITIAL_NODES);
assertEquals(initial, proposed);
// Check against the actual write placements
- assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, initial, INITIAL_NODES);
+ assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), initial, INITIAL_NODES);
// Initiate an operation which affects ownership. This will add the MultiStepOperation which encodes any
// necessary range movements so subsequent calls to ClusterMetadata::localRangesAllSettled
@@ -123,7 +123,7 @@ public class LocalRangesAllSettledTest
assertEquals(proposed, finalized);
// Finally, check against the actual write placements
- assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, finalized, INITIAL_NODES);
+ assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), finalized, INITIAL_NODES);
}
@Test
@@ -134,7 +134,7 @@ public class LocalRangesAllSettledTest
AllLocalRanges proposed = snapshotAllLocalRanges(LocalRangeStatus.SETTLED, INITIAL_NODES);
assertEquals(initial, proposed);
// Check against the actual write placements
- assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, initial, INITIAL_NODES);
+ assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), initial, INITIAL_NODES);
// Initiate an operation which affects ownership. This will add the MultiStepOperation which encodes any
// necessary range movements so subsequent calls to ClusterMetadata::localRangesAllSettled
@@ -161,7 +161,7 @@ public class LocalRangesAllSettledTest
assertEquals(proposed, finalized);
// Finally, check against the actual write placements
- assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, finalized, expandedNodes);
+ assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), finalized, expandedNodes);
}
@Test
@@ -172,7 +172,7 @@ public class LocalRangesAllSettledTest
AllLocalRanges proposed = snapshotAllLocalRanges(LocalRangeStatus.SETTLED, INITIAL_NODES);
assertEquals(initial, proposed);
// Check against the actual write placements
- assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, initial, INITIAL_NODES);
+ assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), initial, INITIAL_NODES);
// Initiate an operation which affects ownership. This will add the MultiStepOperation which encodes any
// necessary range movements so subsequent calls to ClusterMetadata::localRangesAllSettled
@@ -201,7 +201,7 @@ public class LocalRangesAllSettledTest
assertEquals(proposed, finalized);
// Finally, check against the actual write placements
- assertLocalRangesMatchPlacements(ClusterMetadata.current().placements, finalized, INITIAL_NODES);
+ assertLocalRangesMatchPlacements(ClusterMetadata.current().placements(), finalized, INITIAL_NODES);
}
private void assertLocalRangesMatchPlacements(DataPlacements placements,
diff --git a/test/unit/org/apache/cassandra/tcm/ownership/OwnershipUtils.java b/test/unit/org/apache/cassandra/tcm/ownership/OwnershipUtils.java
index 897dea5e5e..4b8b4e4164 100644
--- a/test/unit/org/apache/cassandra/tcm/ownership/OwnershipUtils.java
+++ b/test/unit/org/apache/cassandra/tcm/ownership/OwnershipUtils.java
@@ -254,6 +254,6 @@ public class OwnershipUtils
assert result.isSuccess();
workingMetadata = result.success().metadata;
}
- return workingMetadata.placements;
+ return workingMetadata.placements();
}
}
diff --git a/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java b/test/unit/org/apache/cassandra/tcm/sequences/InProgressSequenceCancellationTest.java
index d4e1b4df29..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
@@ -301,7 +305,7 @@ public class InProgressSequenceCancellationTest
private void assertRelevantMetadata(ClusterMetadata first, ClusterMetadata second)
{
- assertTrue(first.placements.equivalentTo(second.placements));
+ assertTrue(first.placements().equivalentTo(second.placements()));
assertTrue(first.directory.equivalentTo(second.directory));
assertTrue(first.tokenMap.equivalentTo(second.tokenMap));
assertEquals(first.lockedRanges.locked.keySet(), second.lockedRanges.locked.keySet());
diff --git a/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java b/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java
index cb4f15e5ef..2f3a3a0fd0 100644
--- a/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java
+++ b/test/unit/org/apache/cassandra/tcm/sequences/ProgressBarrierTest.java
@@ -113,7 +113,7 @@ public class ProgressBarrierTest extends CMSTestBase
// Internally affectedRanges::toPeers uses the same logic as
// the progress barrier does to identify the consensus group
Set consensusGroup = leave.barrier().affectedRanges.toPeers(ReplicationParams.simple(1),
- sut.service.metadata().placements,
+ sut.service.metadata().placements(),
sut.service.metadata().directory);
assertEquals(Set.of(node2.nodeId(), node3.nodeId()), consensusGroup);
}
@@ -217,7 +217,7 @@ public class ProgressBarrierTest extends CMSTestBase
case ALL:
{
Set replicas = metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
- .toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
+ .toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
.stream()
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
.collect(Collectors.toSet());
@@ -235,7 +235,7 @@ public class ProgressBarrierTest extends CMSTestBase
case QUORUM:
{
Set replicas = metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
- .toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
+ .toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
.stream()
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
.collect(Collectors.toSet());
@@ -253,7 +253,7 @@ public class ProgressBarrierTest extends CMSTestBase
case LOCAL_QUORUM:
{
List replicas = new ArrayList<>(metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
- .toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
+ .toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
.stream()
.filter((n) -> metadata.directory.location(n).datacenter.equals(dc))
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
@@ -277,7 +277,7 @@ public class ProgressBarrierTest extends CMSTestBase
{
Map byDc = new HashMap<>();
metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
- .toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
+ .toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
.forEach(n -> byDc.compute(metadata.directory.location(n).datacenter,
(k, v) -> v == null ? 1 : v + 1));
@@ -306,7 +306,7 @@ public class ProgressBarrierTest extends CMSTestBase
}
case ONE:
Set replicas = metadata.lockedRanges.locked.get(LockedRanges.keyFor(metadata.epoch))
- .toPeers(rf.asKeyspaceParams().replication, metadata.placements, metadata.directory)
+ .toPeers(rf.asKeyspaceParams().replication, metadata.placements(), metadata.directory)
.stream()
.map(n -> metadata.directory.getNodeAddresses(n).broadcastAddress)
.collect(Collectors.toSet());
diff --git a/test/unit/org/apache/cassandra/tcm/transformations/EventsMetadataTest.java b/test/unit/org/apache/cassandra/tcm/transformations/EventsMetadataTest.java
index 11fd03c696..0f9115805f 100644
--- a/test/unit/org/apache/cassandra/tcm/transformations/EventsMetadataTest.java
+++ b/test/unit/org/apache/cassandra/tcm/transformations/EventsMetadataTest.java
@@ -95,8 +95,8 @@ public class EventsMetadataTest
// should not be in tokenMap (no tokens yet)
assertTrue(metadata.tokenMap.tokens(nodeId).isEmpty());
- assertTrue(metadata.placements.get(KSM.params.replication).writes.byEndpoint().isEmpty());
- assertTrue(metadata.placements.get(KSM.params.replication).reads.byEndpoint().isEmpty());
+ assertTrue(metadata.placement(KSM.params.replication).writes.byEndpoint().isEmpty());
+ assertTrue(metadata.placement(KSM.params.replication).reads.byEndpoint().isEmpty());
assertTrue(metadata.lockedRanges.locked.isEmpty());
}
@@ -130,9 +130,9 @@ public class EventsMetadataTest
assertTrue(ClusterMetadata.current().tokenMap.tokens(nodeId).isEmpty());
assertEquals(NodeState.BOOTSTRAPPING, ClusterMetadata.current().directory.peerState(nodeId));
- assertTrue(ClusterMetadata.current().placements.get(KSM.params.replication).writes.byEndpoint().containsKey(node1));
+ assertTrue(ClusterMetadata.current().placement(KSM.params.replication).writes.byEndpoint().containsKey(node1));
// the first joined node gets added to the read endpoints immediately
- assertTrue(ClusterMetadata.current().placements.get(KSM.params.replication).reads.byEndpoint().containsKey(node1));
+ assertTrue(ClusterMetadata.current().placement(KSM.params.replication).reads.byEndpoint().containsKey(node1));
ClusterMetadataService.instance().commit(plan.midJoin);
ClusterMetadataService.instance().commit(plan.finishJoin);
@@ -152,8 +152,8 @@ public class EventsMetadataTest
assertTrue(ClusterMetadata.current().tokenMap.tokens(nodeId).isEmpty());
assertEquals(NodeState.BOOTSTRAPPING, ClusterMetadata.current().directory.peerState(nodeId));
- assertTrue(ClusterMetadata.current().placements.get(KSM.params.replication).writes.byEndpoint().containsKey(node2));
- assertFalse(ClusterMetadata.current().placements.get(KSM.params.replication).reads.byEndpoint().containsKey(node2));
+ assertTrue(ClusterMetadata.current().placement(KSM.params.replication).writes.byEndpoint().containsKey(node2));
+ assertFalse(ClusterMetadata.current().placement(KSM.params.replication).reads.byEndpoint().containsKey(node2));
}
@Test
@@ -178,7 +178,7 @@ public class EventsMetadataTest
// no change in metadata after prepareLeave;
assertEquals(before.directory, after.directory);
assertEquals(before.tokenMap, after.tokenMap);
- assertEquals(before.placements, after.placements);
+ assertEquals(before.placements(), after.placements());
assertEquals(before.schema, after.schema);
ClusterMetadataService.instance().commit(leave.startLeave);
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);
diff --git a/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java b/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java
index 57d87d22e7..04209514bb 100644
--- a/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java
+++ b/test/unit/org/apache/cassandra/tools/CMSOfflineToolTest.java
@@ -22,7 +22,6 @@ import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.ArrayList;
-import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
@@ -55,6 +54,7 @@ import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.service.accord.topology.AccordFastPath;
import org.apache.cassandra.service.accord.topology.AccordStaleReplicas;
import org.apache.cassandra.service.consensus.migration.ConsensusMigrationState;
+import org.apache.cassandra.tcm.CMSMembership;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.tcm.ClusterMetadataService;
import org.apache.cassandra.tcm.Epoch;
@@ -110,7 +110,8 @@ public class CMSOfflineToolTest extends OfflineToolUtils
InProgressSequences.EMPTY,
ConsensusMigrationState.EMPTY,
ImmutableMap.of(),
- AccordStaleReplicas.EMPTY);
+ AccordStaleReplicas.EMPTY,
+ CMSMembership.EMPTY);
}
@Before
@@ -669,7 +670,7 @@ public class CMSOfflineToolTest extends OfflineToolUtils
"Cluster Metadata Service:\n" +
"Members: /127.0.0.1:" + storagePort + ",/127.0.0.2:" + storagePort + ",/127.0.0.3:" + storagePort + '\n' +
"Needs reconfiguration: false\n" +
- "Service State: LOCAL\n" +
+ "Service State: " + ClusterMetadataService.State.OFFLINE_TOOL + '\n' +
"Epoch: 2\n" +
"Replication factor: ReplicationParams{class=org.apache.cassandra.locator.MetaStrategy, datacenter1=3}\n";
assertThat(result.getStdout()).isEqualTo(expectedOutput);
@@ -1668,10 +1669,7 @@ public class CMSOfflineToolTest extends OfflineToolUtils
KeyspaceMetadata normalKeyspace = KeyspaceMetadata.create("ks", KeyspaceParams.simple(3));
Keyspaces keyspaces = Keyspaces.none().with(metaKeyspace).with(normalKeyspace);
- ClusterMetadata clusterMetadata = getClusterMetadata(keyspaces, partitioner, directory);
-
-
- ClusterMetadata metadata = clusterMetadata
+ ClusterMetadata metadata = getClusterMetadata(keyspaces, partitioner, directory)
.transformer()
.with(directory)
.join(nodeId1)
@@ -1680,12 +1678,12 @@ public class CMSOfflineToolTest extends OfflineToolUtils
.proposeToken(nodeId2, getRandomTokens(partitioner, tokenSize))
.join(nodeId3)
.proposeToken(nodeId3, getRandomTokens(partitioner, tokenSize))
+ .startJoiningCMS(nodeId1).finishJoiningCMS(nodeId1)
+ .startJoiningCMS(nodeId2).finishJoiningCMS(nodeId2)
+ .startJoiningCMS(nodeId3).finishJoiningCMS(nodeId3)
.build().metadata;
- // Create replicas for the metadata keyspace on all three nodes
- ReplicationParams metaParams = ReplicationParams.ntsMeta(Collections.singletonMap(DC, 3));
DataPlacements placements = DataPlacements.empty().unbuild()
- .with(metaParams, getCMSMemberPlacement(metadata, List.of(addr1, addr2, addr3)))
.with(ReplicationParams.simple(3), getKeyspacePlacement(metadata, normalKeyspace))
.build();
@@ -1799,6 +1797,11 @@ public class CMSOfflineToolTest extends OfflineToolUtils
ClusterMetadata startReplacing(NodeId oldNodeId, NodeId newNodeId, ClusterMetadata clusterMetadata)
{
+ // In a real cluster, a node being replaced is removed from the CMS prior to the replacement starting
+ // i.e. before the PrepareReplace is committed
+ if (clusterMetadata.fullCMSMemberIds().contains(oldNodeId))
+ clusterMetadata = clusterMetadata.transformer().leaveCMS(oldNodeId).build().metadata;
+
Register register = new Register(getNodeAddresses(newNodeId.id()),
new Location(DC, "rack" + newNodeId.id()),
NodeVersion.CURRENT);
diff --git a/test/unit/org/apache/cassandra/utils/CassandraGenerators.java b/test/unit/org/apache/cassandra/utils/CassandraGenerators.java
index 4c0a6758ec..1438fc8b62 100644
--- a/test/unit/org/apache/cassandra/utils/CassandraGenerators.java
+++ b/test/unit/org/apache/cassandra/utils/CassandraGenerators.java
@@ -140,6 +140,7 @@ import org.apache.cassandra.service.accord.topology.SimpleFastPathStrategy;
import org.apache.cassandra.service.accord.topology.UpFastPathStrategy;
import org.apache.cassandra.service.consensus.TransactionalMode;
import org.apache.cassandra.service.consensus.migration.ConsensusMigrationState;
+import org.apache.cassandra.tcm.CMSMembership;
import org.apache.cassandra.tcm.ClusterMetadata;
import org.apache.cassandra.tcm.Epoch;
import org.apache.cassandra.tcm.extensions.ExtensionKey;
@@ -1955,7 +1956,8 @@ public final class CassandraGenerators
ConsensusMigrationState consensusMigrationState = ConsensusMigrationState.EMPTY;
Map, ExtensionValue>> extensions = ImmutableMap.of();
AccordStaleReplicas accordStaleReplicas = accordStaleReplicasGen.generate(rnd);
- return new ClusterMetadata(epoch, partitioner, schema, directory, tokenMap, placements, accordFastPath, lockedRanges, inProgressSequences, consensusMigrationState, extensions, accordStaleReplicas);
+ CMSMembership cms = CMSMembership.EMPTY;
+ return new ClusterMetadata(epoch, partitioner, schema, directory, tokenMap, placements, accordFastPath, lockedRanges, inProgressSequences, consensusMigrationState, extensions, accordStaleReplicas, cms);
};
}
}