From fe1be800b4f7e6ca5b2f28dddd2b6f7489f41631 Mon Sep 17 00:00:00 2001 From: Abe Ratnofsky Date: Mon, 13 Nov 2023 15:17:47 +0100 Subject: [PATCH] Remove completed coordinator sessions patch by Abe Ratnofsky; reviewed by Caleb Rackliffe, Marcus Eriksson for CASSANDRA-18903 --- CHANGES.txt | 1 + .../repair/consistent/ConsistentSession.java | 6 ++++++ .../repair/consistent/CoordinatorSession.java | 13 +++++++++++++ .../repair/consistent/CoordinatorSessions.java | 14 ++++++++++++++ .../cassandra/repair/consistent/LocalSession.java | 6 ------ .../test/IncRepairCoordinatorErrorTest.java | 2 +- 6 files changed, 35 insertions(+), 7 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 232ac5f56b..6ac97d7a99 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0.12 + * Remove completed coordinator sessions (CASSANDRA-18903) * Make StartupConnectivityChecker only run a connectivity check if there are no nodes which are running a version prior to Cassandra 4 (CASSANDRA-18968) * Retrieve keyspaces metadata and schema version concistently in DescribeStatement (CASSANDRA-18921) * Gossip NPE due to shutdown event corrupting empty statuses (CASSANDRA-18913) diff --git a/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java b/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java index e4d8ff019a..3d6ace39e3 100644 --- a/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java +++ b/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java @@ -206,6 +206,12 @@ public abstract class ConsistentSession this.participants = ImmutableSet.copyOf(builder.participants); } + public boolean isCompleted() + { + State s = getState(); + return s == State.FINALIZED || s == State.FAILED; + } + public State getState() { return state; diff --git a/src/java/org/apache/cassandra/repair/consistent/CoordinatorSession.java b/src/java/org/apache/cassandra/repair/consistent/CoordinatorSession.java index 5ddac3f745..ef49fa7ed9 100644 --- a/src/java/org/apache/cassandra/repair/consistent/CoordinatorSession.java +++ b/src/java/org/apache/cassandra/repair/consistent/CoordinatorSession.java @@ -23,6 +23,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.function.Consumer; import java.util.function.Supplier; import java.util.stream.Collectors; @@ -71,9 +72,12 @@ public class CoordinatorSession extends ConsistentSession private volatile long repairStart = Long.MIN_VALUE; private volatile long finalizeStart = Long.MIN_VALUE; + private final Consumer listener; + public CoordinatorSession(Builder builder) { super(builder); + this.listener = builder.listener; for (InetAddressAndPort participant : participants) { participantStates.put(participant, State.PREPARING); @@ -82,6 +86,13 @@ public class CoordinatorSession extends ConsistentSession public static class Builder extends AbstractBuilder { + Consumer listener; + + public void withListener(Consumer listener) + { + this.listener = listener; + } + public CoordinatorSession build() { validate(); @@ -98,6 +109,8 @@ public class CoordinatorSession extends ConsistentSession { logger.trace("Setting coordinator state to {} for repair {}", state, sessionID); super.setState(state); + if (listener != null) + listener.accept(this); } @VisibleForTesting diff --git a/src/java/org/apache/cassandra/repair/consistent/CoordinatorSessions.java b/src/java/org/apache/cassandra/repair/consistent/CoordinatorSessions.java index b87a2c085c..f14d92ecc6 100644 --- a/src/java/org/apache/cassandra/repair/consistent/CoordinatorSessions.java +++ b/src/java/org/apache/cassandra/repair/consistent/CoordinatorSessions.java @@ -25,6 +25,9 @@ import java.util.UUID; import com.google.common.base.Preconditions; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.repair.messages.FailSession; import org.apache.cassandra.repair.messages.FinalizePromise; @@ -37,6 +40,7 @@ import org.apache.cassandra.service.ActiveRepairService; public class CoordinatorSessions { private final Map sessions = new HashMap<>(); + private static final Logger logger = LoggerFactory.getLogger(CoordinatorSessions.class); protected CoordinatorSession buildSession(CoordinatorSession.Builder builder) { @@ -61,6 +65,7 @@ public class CoordinatorSessions builder.withRepairedAt(prs.repairedAt); builder.withRanges(prs.getRanges()); builder.withParticipants(participants); + builder.withListener(this::onSessionStateUpdate); CoordinatorSession session = buildSession(builder); sessions.put(session.sessionID, session); return session; @@ -71,6 +76,15 @@ public class CoordinatorSessions return sessions.get(sessionId); } + public synchronized void onSessionStateUpdate(CoordinatorSession session) + { + if (session.isCompleted()) + { + logger.info("Removing completed session {} with state {}", session.sessionID, session.getState()); + sessions.remove(session.sessionID); + } + } + public void handlePrepareResponse(PrepareConsistentResponse msg) { CoordinatorSession session = getSession(msg.parentSession); diff --git a/src/java/org/apache/cassandra/repair/consistent/LocalSession.java b/src/java/org/apache/cassandra/repair/consistent/LocalSession.java index a6f81d7fe5..204dd2a8b5 100644 --- a/src/java/org/apache/cassandra/repair/consistent/LocalSession.java +++ b/src/java/org/apache/cassandra/repair/consistent/LocalSession.java @@ -37,12 +37,6 @@ public class LocalSession extends ConsistentSession this.lastUpdate = builder.lastUpdate; } - public boolean isCompleted() - { - State s = getState(); - return s == State.FINALIZED || s == State.FAILED; - } - public int getStartedAt() { return startedAt; diff --git a/test/distributed/org/apache/cassandra/distributed/test/IncRepairCoordinatorErrorTest.java b/test/distributed/org/apache/cassandra/distributed/test/IncRepairCoordinatorErrorTest.java index c06e848399..148cbd0d03 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/IncRepairCoordinatorErrorTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/IncRepairCoordinatorErrorTest.java @@ -51,7 +51,7 @@ public class IncRepairCoordinatorErrorTest extends TestBaseImpl cluster.get(3).runOnInstance(() -> { ActiveRepairService.instance.failSession(result.toString(), true); }); - assertThat(cluster.get(1).logs().watchFor("Can't transition endpoints .* to FAILED").getResult()).isNotEmpty(); + assertThat(cluster.get(1).logs().watchFor("Removing completed session .* with state FINALIZED").getResult()).isNotEmpty(); } } }