diff --git a/src/java/org/apache/cassandra/tcm/Commit.java b/src/java/org/apache/cassandra/tcm/Commit.java index 847cd22a7ca3..dab4f05cdef6 100644 --- a/src/java/org/apache/cassandra/tcm/Commit.java +++ b/src/java/org/apache/cassandra/tcm/Commit.java @@ -48,6 +48,7 @@ import org.apache.cassandra.tcm.membership.Directory; import org.apache.cassandra.tcm.membership.EndpointLookup; import org.apache.cassandra.tcm.membership.NodeId; +import org.apache.cassandra.tcm.membership.NodeState; import org.apache.cassandra.tcm.membership.NodeVersion; import org.apache.cassandra.tcm.serialization.Version; import org.apache.cassandra.utils.FBUtilities; @@ -461,6 +462,11 @@ InetAddressAndPort endpoint(NodeId id) { return endpoints.endpoint(id); }; + + boolean hasLeft(NodeId id) + { + return directory.peerState(id) == NodeState.LEFT; + } } private final Supplier routingSupplier; @@ -490,10 +496,12 @@ public void send(Result result, InetAddressAndPort source) { InetAddressAndPort endpoint = routing.endpoint(peerId); boolean upgraded = routing.isUpgraded(peerId); - // Do not replicate to self and to the peer that has requested to commit this message + boolean hasLeft = routing.hasLeft(peerId); + // Do not replicate to self, to any left nodes and to the peer that has requested to commit this message if (endpoint.equals(FBUtilities.getBroadcastAddressAndPort()) || (source != null && source.equals(endpoint)) || - !upgraded) + !upgraded || + hasLeft) { continue; } diff --git a/test/distributed/org/apache/cassandra/distributed/test/ring/DecommissionTest.java b/test/distributed/org/apache/cassandra/distributed/test/ring/DecommissionTest.java index 76f53e9e6e87..864be8ad2248 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/ring/DecommissionTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/ring/DecommissionTest.java @@ -27,6 +27,7 @@ import java.util.concurrent.Callable; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import java.util.stream.Collectors; @@ -438,5 +439,20 @@ public void testPeersPostDecom() throws IOException } } + @Test + public void testDontReplicateToLeftNodes() throws IOException, TimeoutException + { + try (Cluster cluster = init(builder().withNodes(3) + .withConfig(config -> config.with(NETWORK, GOSSIP)) + .start())) + { + cluster.get(2).nodetoolResult("decommission", "--force").asserts().success(); + long mark = cluster.get(1).logs().mark(); + cluster.coordinator(1).execute(withKeyspace("create table %s.tbl (id int primary key)"), ConsistencyLevel.ONE); + cluster.get(3).logs().watchFor("Enacted.*CreateTableStatement.*tbl"); + assertTrue(cluster.get(1).logs().grep(mark, "Replicating newly committed transformations up to.*127.0.0.2.*").getResult().isEmpty()); + } + } + } diff --git a/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java b/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java index e8429e5bb723..f2c38d619d23 100644 --- a/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java +++ b/test/unit/org/apache/cassandra/tcm/listeners/PlacementsChangeListenerTest.java @@ -39,7 +39,6 @@ import org.apache.cassandra.tcm.ClusterMetadata; import org.apache.cassandra.tcm.Epoch; import org.apache.cassandra.tcm.membership.Directory; -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.ownership.DataPlacement; @@ -81,7 +80,8 @@ public void testPlacementChange() 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()); + // randomPlacements gives 127.0.0.{1-255} ip addresses, need to add a non-conflicting one here: + Replica add = Replica.fullReplica(InetAddressAndPort.getByNameUnchecked("127.0.1.99"), remove.range()); DataPlacement newPlacement = placement.unbuild() .withoutWriteReplica(e.nextEpoch(), remove) .withWriteReplica(e.nextEpoch(), add).build();