diff --git a/CHANGES.txt b/CHANGES.txt index 546e3449d3..44cc5dcce5 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -34,6 +34,7 @@ Merged from 5.0: * Use estimated compressed size for tables to check if there is enough free space for a compaction (CASSANDRA-21245) * Fix failing select on system_views.settings for non-string keys (CASSANDRA-21348) Merged from 4.0: + * Remove inFlightEcho entry on ECHO_REQ failure (CASSANDRA-21428) * Validate snapshot names (CASSANDRA-21389) diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index e35ba861af..979a79a4b2 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -69,6 +69,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Token; +import org.apache.cassandra.exceptions.RequestFailure; import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -659,6 +660,7 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean, public void evictFromMembership(InetAddressAndPort endpoint) { checkProperThreadForStateMutation(); + inflightEcho.remove(endpoint); unreachableEndpoints.remove(endpoint); endpointStateMap.remove(endpoint); expireTimeEndpointMap.remove(endpoint); @@ -1201,13 +1203,30 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean, { Message echoMessage = Message.out(ECHO_REQ, noPayload); logger.trace("Sending ECHO_REQ to {}", addr); - RequestCallback echoHandler = msg -> + RequestCallback echoHandler = new RequestCallback() { - runInGossipStageBlocking(() -> { - EndpointState epState = inflightEcho.remove(addr); - if (epState != null) - realMarkAlive(addr, epState); - }); + @Override + public void onResponse(Message msg) + { + runInGossipStageBlocking(() -> { + EndpointState epState = inflightEcho.remove(addr); + if (epState != null) + realMarkAlive(addr, epState); + }); + } + + @Override + public boolean invokeOnFailure() + { + return true; + } + + @Override + public void onFailure(InetAddressAndPort from, RequestFailure requestFailure) + { + logger.trace("ECHO_REQ to {} failed ({})", addr, requestFailure); + inflightEcho.remove(addr); + } }; MessagingService.instance().sendWithCallback(echoMessage, addr, echoHandler); }