diff --git a/CHANGES.txt b/CHANGES.txt index e27f07de82..5a4431901a 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -46,6 +46,7 @@ Merged from 4.1: * Add Paxos v2 option and informatin in cassandra.yaml (CASSANDRA-21316) * Harden data resurrection startup check with atomic heartbeat file write with fallback (CASSANDRA-21290) Merged from 4.0: + * Remove inFlightEcho entry on ECHO_REQ failure (CASSANDRA-21428) * Validate snapshot names (CASSANDRA-21389) * BTree.FastBuilder.reset() fails to clear savedBuffer and savedNextKey, causing ClassCastException and SSTable header corruption during schema disagreement (CASSANDRA-21216, CASSANDRA-21260) * Backport CASSANDRA-17810 fix and improve RTBoundValidator error messages (CASSANDRA-18282) diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index e35ba861af..728ae7de7c 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 failure) + { + logger.trace("ECHO_REQ to {} failed ({})", addr, failure); + inflightEcho.remove(addr); + } }; MessagingService.instance().sendWithCallback(echoMessage, addr, echoHandler); }