Merge branch 'cassandra-4.1' into cassandra-5.0

This commit is contained in:
Stefan Miklosovic 2023-11-13 15:24:17 +01:00
commit bd4e7d7824
No known key found for this signature in database
GPG Key ID: 32F35CB2F546D93E
6 changed files with 36 additions and 7 deletions

View File

@ -5,6 +5,8 @@
* Add metrics and logging to repair retries (CASSANDRA-18952)
* Remove deprecated code in Cassandra 1.x and 2.x (CASSANDRA-18959)
* ClientRequestSize metrics should not treat CONTAINS restrictions as being equality-based (CASSANDRA-18896)
Merged from 4.0:
* Remove completed coordinator sessions (CASSANDRA-18903)
5.0-alpha2

View File

@ -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;

View File

@ -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<CoordinatorSession> 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<CoordinatorSession> listener;
private SharedContext ctx;
public Builder(SharedContext ctx)
@ -94,6 +99,11 @@ public class CoordinatorSession extends ConsistentSession
super(ctx);
}
public void withListener(Consumer<CoordinatorSession> 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

View File

@ -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<TimeUUID, CoordinatorSession> 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<? extends RepairMessage> msg)
{
PrepareConsistentResponse payload = (PrepareConsistentResponse) msg.payload;

View File

@ -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;

View File

@ -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();
}
}
}