diff --git a/CHANGES.txt b/CHANGES.txt index a2c41d017a..c77d936885 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -37,6 +37,7 @@ Merged from 4.1: * Internode legacy SSL storage port certificate is not hot reloaded on update (CASSANDRA-18681) * Nodetool paxos-only repair is no longer incremental (CASSANDRA-18466) Merged from 4.0: + * Remove completed coordinator sessions (CASSANDRA-18903) * Gossip NPE due to shutdown event corrupting empty statuses (CASSANDRA-18913) * Update hdrhistogram to 2.1.12 (CASSANDRA-18893) * Improve performance of compactions when table does not have an index (CASSANDRA-18773) diff --git a/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java b/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java index 8101f02912..d5eb5edb31 100644 --- a/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java +++ b/src/java/org/apache/cassandra/repair/consistent/ConsistentSession.java @@ -209,6 +209,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 6771dfe07d..d91960adf8 100644 --- a/src/java/org/apache/cassandra/repair/consistent/CoordinatorSession.java +++ b/src/java/org/apache/cassandra/repair/consistent/CoordinatorSession.java @@ -21,6 +21,7 @@ package org.apache.cassandra.repair.consistent; import java.util.HashMap; import java.util.Map; import java.util.Set; +import java.util.function.Consumer; import java.util.function.Supplier; import java.util.stream.Collectors; @@ -75,9 +76,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; ctx = builder.ctx == null ? SharedContext.Global.instance : builder.ctx; for (InetAddressAndPort participant : participants) { @@ -87,6 +91,7 @@ public class CoordinatorSession extends ConsistentSession public static class Builder extends AbstractBuilder { + Consumer listener; private SharedContext ctx; public Builder(SharedContext ctx) @@ -94,6 +99,11 @@ public class CoordinatorSession extends ConsistentSession super(ctx); } + public void withListener(Consumer listener) + { + this.listener = listener; + } + public Builder withContext(SharedContext ctx) { this.ctx = ctx; @@ -116,6 +126,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 6f02f033d3..407ad9f788 100644 --- a/src/java/org/apache/cassandra/repair/consistent/CoordinatorSessions.java +++ b/src/java/org/apache/cassandra/repair/consistent/CoordinatorSessions.java @@ -24,6 +24,9 @@ import java.util.Set; import com.google.common.base.Preconditions; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import org.apache.cassandra.net.Message; import org.apache.cassandra.repair.SharedContext; import org.apache.cassandra.locator.InetAddressAndPort; @@ -42,6 +45,7 @@ import static org.apache.cassandra.repair.messages.RepairMessage.sendFailureResp */ public class CoordinatorSessions { + private static final Logger logger = LoggerFactory.getLogger(CoordinatorSessions.class); private final SharedContext ctx; private final Map sessions = new HashMap<>(); @@ -74,6 +78,7 @@ public class CoordinatorSessions builder.withRepairedAt(prs.repairedAt); builder.withRanges(prs.getRanges()); builder.withParticipants(participants); + builder.withListener(this::onSessionStateUpdate); builder.withContext(ctx); CoordinatorSession session = buildSession(builder); sessions.put(session.sessionID, session); @@ -85,6 +90,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(Message msg) { PrepareConsistentResponse payload = (PrepareConsistentResponse) msg.payload; diff --git a/src/java/org/apache/cassandra/repair/consistent/LocalSession.java b/src/java/org/apache/cassandra/repair/consistent/LocalSession.java index 06c6755562..a7ddea6ba9 100644 --- a/src/java/org/apache/cassandra/repair/consistent/LocalSession.java +++ b/src/java/org/apache/cassandra/repair/consistent/LocalSession.java @@ -39,12 +39,6 @@ public class LocalSession extends ConsistentSession this.lastUpdate = builder.lastUpdate; } - public boolean isCompleted() - { - State s = getState(); - return s == State.FINALIZED || s == State.FAILED; - } - public long 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 06ac200be5..35a863417e 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/IncRepairCoordinatorErrorTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/IncRepairCoordinatorErrorTest.java @@ -46,11 +46,12 @@ public class IncRepairCoordinatorErrorTest extends TestBaseImpl .to(3) .messagesMatching((from, to, msg) -> msg.verb() == FINALIZE_COMMIT_MSG.id).drop(); cluster.get(1).nodetoolResult("repair", KEYSPACE).asserts().success(); + assertThat(cluster.get(1).logs().watchFor("Removing completed session .* with state FINALIZED").getResult()).isNotEmpty(); + TimeUUID result = (TimeUUID) cluster.get(1).executeInternal("select parent_id from system_distributed.repair_history")[0][0]; cluster.get(3).runOnInstance(() -> { ActiveRepairService.instance().failSession(result.toString(), true); }); - assertThat(cluster.get(1).logs().watchFor("Can't transition endpoints .* to FAILED").getResult()).isNotEmpty(); } } }