From 32755cabfa2eeb99f0b8c91fc7bb53379259de54 Mon Sep 17 00:00:00 2001 From: Sam Tunnicliffe Date: Thu, 18 Jul 2024 19:23:46 +0100 Subject: [PATCH] Correctly update peers tables following replacement Patch by Sam Tunnicliffe; reviewed by David Capwell for CASSANDRA-19782 --- CHANGES.txt | 1 + .../tcm/listeners/LegacyStateListener.java | 1 + .../hostreplacement/HostReplacementTest.java | 36 +++++++++++++++++++ 3 files changed, 38 insertions(+) diff --git a/CHANGES.txt b/CHANGES.txt index a8f7d35c2c..24b68b6a03 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.1 + * Update legacy peers tables during node replacement (CASSANDRA-19782) * Refactor ColumnCondition (CASSANDRA-19620) * Allow configuring log format for Audit Logs (CASSANDRA-19792) * Support for noboolean rpm (centos7 compatible) packages removed (CASSANDRA-19787) diff --git a/src/java/org/apache/cassandra/tcm/listeners/LegacyStateListener.java b/src/java/org/apache/cassandra/tcm/listeners/LegacyStateListener.java index 3a8701f8f2..25aca06f00 100644 --- a/src/java/org/apache/cassandra/tcm/listeners/LegacyStateListener.java +++ b/src/java/org/apache/cassandra/tcm/listeners/LegacyStateListener.java @@ -158,6 +158,7 @@ public class LegacyStateListener implements ChangeListener.Async logger.warn("Token {} changing ownership from {} to {}", token, replaced, replacement); } Gossiper.instance.mergeNodeToGossip(change, next, tokens); + PeersTable.updateLegacyPeerTable(change, prev, next); } } else diff --git a/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/HostReplacementTest.java b/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/HostReplacementTest.java index 85105dc8a3..34520636c6 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/HostReplacementTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/hostreplacement/HostReplacementTest.java @@ -19,8 +19,11 @@ package org.apache.cassandra.distributed.test.hostreplacement; import java.io.IOException; +import java.net.InetAddress; import java.util.Arrays; +import java.util.LinkedHashSet; import java.util.List; +import java.util.UUID; import java.util.concurrent.TimeUnit; import org.junit.Test; @@ -41,11 +44,15 @@ import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.schema.SchemaConstants; +import org.apache.cassandra.tcm.ClusterMetadata; +import org.apache.cassandra.tcm.membership.NodeId; +import org.apache.cassandra.tcm.membership.NodeState; import org.assertj.core.api.Assertions; import static org.apache.cassandra.config.CassandraRelevantProperties.BOOTSTRAP_SKIP_SCHEMA_CHECK; import static org.apache.cassandra.config.CassandraRelevantProperties.GOSSIPER_QUARANTINE_DELAY; import static org.apache.cassandra.distributed.shared.AssertUtils.assertRows; +import static org.apache.cassandra.distributed.shared.AssertUtils.row; import static org.apache.cassandra.distributed.shared.ClusterUtils.assertInRing; import static org.apache.cassandra.distributed.shared.ClusterUtils.assertRingIs; import static org.apache.cassandra.distributed.shared.ClusterUtils.awaitRingHealthy; @@ -115,6 +122,7 @@ public class HostReplacementTest extends TestBaseImpl String replacingNodeAddress = replacingNode.config().broadcastAddress().getHostString(); validateGossipStatusNormal(seed, replacingNodeAddress); validateGossipStatusNormal(replacingNode, replacingNodeAddress); + validatePeersTables(seed, replacingNode, even); } } @@ -266,4 +274,32 @@ public class HostReplacementTest extends TestBaseImpl }); } + static void validatePeersTables(IInvokableInstance seed, IInvokableInstance replacingNode, TokenSupplier tokens) + { + InetAddress replacementAddress = replacingNode.config().broadcastAddress().getAddress(); + NodeId replacementId = new NodeId(3); + UUID schemaVersion = seed.callOnInstance(() -> ClusterMetadata.current().schema.getVersion()); + String releaseVersion = seed.callOnInstance(() -> ClusterMetadata.current().directory.version(ClusterMetadata.current().myNodeId()).cassandraVersion.toString()); + LinkedHashSet replacementTokens = new LinkedHashSet<>(tokens.tokens(2)); + String datacenter = "datacenter0"; + String rack = "rack0"; + assertRows(seed.executeInternal("SELECT * from system.peers"), + row(replacementAddress, datacenter, replacementId.toUUID(), + replacementAddress, rack, releaseVersion, replacementAddress, + schemaVersion, replacementTokens)); + assertRows(seed.executeInternal("SELECT * from system.peers_v2"), + row(replacementAddress, 7012, datacenter, replacementId.toUUID(), + replacementAddress, 9042, replacementAddress, 7012, rack, releaseVersion, + schemaVersion, replacementTokens)); + // system_views.peers contains both remote and local node info + InetAddress seedAddress = seed.config().broadcastAddress().getAddress(); + assertRows(seed.executeInternal("SELECT * from system_views.peers"), + rows(row(seedAddress, 7012, datacenter, new NodeId(1).toUUID(), + seedAddress, 9042, seedAddress, 7012, rack, releaseVersion, + schemaVersion, NodeState.JOINED.toString(), new LinkedHashSet<>(tokens.tokens(1))), + row(replacementAddress, 7012, datacenter, replacementId.toUUID(), + replacementAddress, 9042, replacementAddress, 7012, rack, releaseVersion, + schemaVersion, NodeState.JOINED.toString(), replacementTokens))); + + } }