From 3ad4794aa006c239fb60923449d96bb8ae70d891 Mon Sep 17 00:00:00 2001 From: Stefan Miklosovic Date: Wed, 4 Mar 2026 09:07:17 +0100 Subject: [PATCH] Node does not send multiple inflight echos patch by Cameron Zemek; reviewed by Stefan Miklosovic, Brandon Williams for CASSANDRA-18866 --- CHANGES.txt | 1 + .../org/apache/cassandra/gms/Gossiper.java | 30 ++++++++++++++----- 2 files changed, 24 insertions(+), 7 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 271fb4ec6e..6c5b443133 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0.20 + * Node does not send multiple inflight echos (CASSANDRA-18866) * Obsolete expired SSTables before compaction starts (CASSANDRA-19776) * No need to evict already prepared statements, as it creates a race condition between multiple threads (CASSANDRA-17401) * Switch lz4-java to at.yawk.lz4 version due to CVE (CASSANDRA-21052) diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index 12c532b162..867bcbbca7 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -140,6 +140,9 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean @VisibleForTesting final Set liveEndpoints = new ConcurrentSkipListSet<>(); + /* Inflight echo requests. */ + private final Map inflightEcho = new ConcurrentHashMap<>(); + /* unreachable member set */ private final Map unreachableEndpoints = new ConcurrentHashMap<>(); @@ -180,6 +183,7 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean { unreachableEndpoints.clear(); liveEndpoints.clear(); + inflightEcho.clear(); justRemovedEndpoints.clear(); expireTimeEndpointMap.clear(); endpointStateMap.clear(); @@ -659,6 +663,7 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean } liveEndpoints.remove(endpoint); + inflightEcho.remove(endpoint); unreachableEndpoints.remove(endpoint); MessagingService.instance().versions.reset(endpoint); quarantineEndpoint(endpoint); @@ -1328,14 +1333,24 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean { localState.markDead(); - Message echoMessage = Message.out(ECHO_REQ, noPayload); - logger.trace("Sending ECHO_REQ to {}", addr); - RequestCallback echoHandler = msg -> + EndpointState prevState = inflightEcho.put(addr, localState); + boolean sendEcho = !localState.equals(prevState); + if (sendEcho) { - runInGossipStageBlocking(() -> realMarkAlive(addr, localState)); - }; - - MessagingService.instance().sendWithCallback(echoMessage, addr, echoHandler); + Message echoMessage = Message.out(ECHO_REQ, noPayload); + logger.trace("Sending ECHO_REQ to {}", addr); + RequestCallback echoHandler = msg -> + { + runInGossipStageBlocking(() -> { + EndpointState epState = inflightEcho.remove(addr); + if (epState != null) + realMarkAlive(addr, epState); + }); + }; + MessagingService.instance().sendWithCallback(echoMessage, addr, echoHandler); + } + else + logger.trace("Skipping ECHO_REQ to {} since it is already inflight", addr); GossiperDiagnostics.markedAlive(this, addr, localState); } @@ -1386,6 +1401,7 @@ public class Gossiper implements IFailureDetectionEventListener, GossiperMBean { localState.markDead(); liveEndpoints.remove(addr); + inflightEcho.remove(addr); unreachableEndpoints.put(addr, System.nanoTime()); }