diff --git a/CHANGES.txt b/CHANGES.txt index d9be2ff393..6320d4e5cb 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -2,6 +2,7 @@ * 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 d443658f83..a01dc0ba7a 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -63,6 +63,7 @@ import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.SystemKeyspace; import org.apache.cassandra.dht.Token; +import org.apache.cassandra.exceptions.RequestFailureReason; import org.apache.cassandra.net.RequestCallback; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -656,6 +657,7 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean private void evictFromMembership(InetAddressAndPort endpoint) { checkProperThreadForStateMutation(); + inflightEcho.remove(endpoint); unreachableEndpoints.remove(endpoint); endpointStateMap.remove(endpoint); expireTimeEndpointMap.remove(endpoint); @@ -1396,13 +1398,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, RequestFailureReason failureReason) + { + logger.trace("ECHO_REQ to {} failed ({})", addr, failureReason); + inflightEcho.remove(addr); + } }; MessagingService.instance().sendWithCallback(echoMessage, addr, echoHandler); }