Skip to content
Merged
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import org.apache.kafka.raft.MockLog.LogBatch;
import org.apache.kafka.raft.MockLog.LogEntry;
import org.apache.kafka.raft.internals.BatchMemoryPool;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;

import java.io.IOException;
Expand Down Expand Up @@ -59,6 +60,7 @@
import static org.junit.jupiter.api.Assertions.fail;
import static org.junit.jupiter.api.Assumptions.assumeTrue;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is assumeTrue used in this test? I test ```RaftEventSimulationTest`` but there is no ignored test cases.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, removed since it's not currently needed.


@Tag("integration")
public class RaftEventSimulationTest {
private static final TopicPartition METADATA_PARTITION = new TopicPartition("__cluster_metadata", 0);
private static final int ELECTION_TIMEOUT_MS = 1000;
Expand Down Expand Up @@ -140,44 +142,44 @@ public void testElectionAfterLeaderFailureQuorumSizeFiveAndThreeObservers() {
}

private void testElectionAfterLeaderFailure(QuorumConfig config) {
testElectionAfterLeaderShutdown(config, false);
checkElectionAfterLeaderShutdown(config, false);
}

@Test
public void testElectionAfterLeaderGracefulShutdownQuorumSizeThree() {
testElectionAfterLeaderGracefulShutdown(new QuorumConfig(3, 0));
checkElectionAfterLeaderGracefulShutdown(new QuorumConfig(3, 0));
}

@Test
public void testElectionAfterLeaderGracefulShutdownQuorumSizeThreeAndTwoObservers() {
testElectionAfterLeaderGracefulShutdown(new QuorumConfig(3, 2));
checkElectionAfterLeaderGracefulShutdown(new QuorumConfig(3, 2));
}

@Test
public void testElectionAfterLeaderGracefulShutdownQuorumSizeFour() {
testElectionAfterLeaderGracefulShutdown(new QuorumConfig(4, 0));
checkElectionAfterLeaderGracefulShutdown(new QuorumConfig(4, 0));
}

@Test
public void testElectionAfterLeaderGracefulShutdownQuorumSizeFourAndTwoObservers() {
testElectionAfterLeaderGracefulShutdown(new QuorumConfig(4, 2));
checkElectionAfterLeaderGracefulShutdown(new QuorumConfig(4, 2));
}

@Test
public void testElectionAfterLeaderGracefulShutdownQuorumSizeFive() {
testElectionAfterLeaderGracefulShutdown(new QuorumConfig(5, 0));
checkElectionAfterLeaderGracefulShutdown(new QuorumConfig(5, 0));
}

@Test
public void testElectionAfterLeaderGracefulShutdownQuorumSizeFiveAndThreeObservers() {
testElectionAfterLeaderGracefulShutdown(new QuorumConfig(5, 3));
checkElectionAfterLeaderGracefulShutdown(new QuorumConfig(5, 3));
}

private void testElectionAfterLeaderGracefulShutdown(QuorumConfig config) {
testElectionAfterLeaderShutdown(config, true);
private void checkElectionAfterLeaderGracefulShutdown(QuorumConfig config) {
checkElectionAfterLeaderShutdown(config, true);
}

private void testElectionAfterLeaderShutdown(QuorumConfig config, boolean isGracefulShutdown) {
private void checkElectionAfterLeaderShutdown(QuorumConfig config, boolean isGracefulShutdown) {
// We need at least three voters to run this tests
assumeTrue(config.numVoters > 2);

Expand Down Expand Up @@ -214,20 +216,20 @@ private void testElectionAfterLeaderShutdown(QuorumConfig config, boolean isGrac

@Test
public void testRecoveryAfterAllNodesFailQuorumSizeThree() {
testRecoveryAfterAllNodesFail(new QuorumConfig(3));
checkRecoveryAfterAllNodesFail(new QuorumConfig(3));
}

@Test
public void testRecoveryAfterAllNodesFailQuorumSizeFour() {
testRecoveryAfterAllNodesFail(new QuorumConfig(4));
checkRecoveryAfterAllNodesFail(new QuorumConfig(4));
}

@Test
public void testRecoveryAfterAllNodesFailQuorumSizeFive() {
testRecoveryAfterAllNodesFail(new QuorumConfig(5));
checkRecoveryAfterAllNodesFail(new QuorumConfig(5));
}

private void testRecoveryAfterAllNodesFail(QuorumConfig config) {
private void checkRecoveryAfterAllNodesFail(QuorumConfig config) {
for (int seed = 0; seed < 100; seed++) {
Cluster cluster = new Cluster(config, seed);
MessageRouter router = new MessageRouter(cluster);
Expand Down Expand Up @@ -259,35 +261,35 @@ private void testRecoveryAfterAllNodesFail(QuorumConfig config) {

@Test
public void testElectionAfterLeaderNetworkPartitionQuorumSizeThree() {
testElectionAfterLeaderNetworkPartition(new QuorumConfig(3));
checkElectionAfterLeaderNetworkPartition(new QuorumConfig(3));
}

@Test
public void testElectionAfterLeaderNetworkPartitionQuorumSizeThreeAndTwoObservers() {
testElectionAfterLeaderNetworkPartition(new QuorumConfig(3, 2));
checkElectionAfterLeaderNetworkPartition(new QuorumConfig(3, 2));
}

@Test
public void testElectionAfterLeaderNetworkPartitionQuorumSizeFour() {
testElectionAfterLeaderNetworkPartition(new QuorumConfig(4));
checkElectionAfterLeaderNetworkPartition(new QuorumConfig(4));
}

@Test
public void testElectionAfterLeaderNetworkPartitionQuorumSizeFourAndTwoObservers() {
testElectionAfterLeaderNetworkPartition(new QuorumConfig(4, 2));
checkElectionAfterLeaderNetworkPartition(new QuorumConfig(4, 2));
}

@Test
public void testElectionAfterLeaderNetworkPartitionQuorumSizeFive() {
testElectionAfterLeaderNetworkPartition(new QuorumConfig(5));
checkElectionAfterLeaderNetworkPartition(new QuorumConfig(5));
}

@Test
public void testElectionAfterLeaderNetworkPartitionQuorumSizeFiveAndThreeObservers() {
testElectionAfterLeaderNetworkPartition(new QuorumConfig(5, 3));
checkElectionAfterLeaderNetworkPartition(new QuorumConfig(5, 3));
}

private void testElectionAfterLeaderNetworkPartition(QuorumConfig config) {
private void checkElectionAfterLeaderNetworkPartition(QuorumConfig config) {
// We need at least three voters to run this tests
assumeTrue(config.numVoters > 2);

Expand Down Expand Up @@ -318,15 +320,15 @@ private void testElectionAfterLeaderNetworkPartition(QuorumConfig config) {

@Test
public void testElectionAfterMultiNodeNetworkPartitionQuorumSizeFive() {
testElectionAfterMultiNodeNetworkPartition(new QuorumConfig(5));
checkElectionAfterMultiNodeNetworkPartition(new QuorumConfig(5));
}

@Test
public void testElectionAfterMultiNodeNetworkPartitionQuorumSizeFiveAndTwoObservers() {
testElectionAfterMultiNodeNetworkPartition(new QuorumConfig(5, 2));
checkElectionAfterMultiNodeNetworkPartition(new QuorumConfig(5, 2));
}

private void testElectionAfterMultiNodeNetworkPartition(QuorumConfig config) {
private void checkElectionAfterMultiNodeNetworkPartition(QuorumConfig config) {
// We need at least three voters to run this tests
assumeTrue(config.numVoters > 2);

Expand Down Expand Up @@ -372,15 +374,15 @@ private void testElectionAfterMultiNodeNetworkPartition(QuorumConfig config) {

@Test
public void testBackToBackLeaderFailuresQuorumSizeThree() {
testBackToBackLeaderFailures(new QuorumConfig(3));
checkBackToBackLeaderFailures(new QuorumConfig(3));
}

@Test
public void testBackToBackLeaderFailuresQuorumSizeFiveAndTwoObservers() {
testBackToBackLeaderFailures(new QuorumConfig(5, 2));
checkBackToBackLeaderFailures(new QuorumConfig(5, 2));
}

private void testBackToBackLeaderFailures(QuorumConfig config) {
private void checkBackToBackLeaderFailures(QuorumConfig config) {
for (int seed = 0; seed < 100; seed++) {
Cluster cluster = new Cluster(config, seed);
MessageRouter router = new MessageRouter(cluster);
Expand Down Expand Up @@ -430,24 +432,19 @@ private void schedulePolling(EventScheduler scheduler,
}
}

@FunctionalInterface
private interface Action {
void execute();
}

private static abstract class Event implements Comparable<Event> {
final int eventId;
final long deadlineMs;
final Action action;
final Runnable action;

protected Event(Action action, int eventId, long deadlineMs) {
protected Event(Runnable action, int eventId, long deadlineMs) {
this.action = action;
this.eventId = eventId;
this.deadlineMs = deadlineMs;
}

void execute(EventScheduler scheduler) {
action.execute();
action.run();
}

public int compareTo(Event other) {
Expand All @@ -463,7 +460,7 @@ private static class PeriodicEvent extends Event {
final int periodMs;
final int jitterMs;

protected PeriodicEvent(Action action,
protected PeriodicEvent(Runnable action,
int eventId,
Random random,
long deadlineMs,
Expand All @@ -483,15 +480,15 @@ void execute(EventScheduler scheduler) {
}
}

private static class SequentialAppendAction implements Action {
private static class SequentialAppendAction implements Runnable {
final Cluster cluster;

private SequentialAppendAction(Cluster cluster) {
this.cluster = cluster;
}

@Override
public void execute() {
public void run() {
cluster.withCurrentLeader(node -> {
if (!node.client.isShuttingDown() && node.counter.isWritable())
node.counter.increment();
Expand Down Expand Up @@ -531,7 +528,7 @@ private void addValidation(Validation validation) {
validations.add(validation);
}

void schedule(Action action, int delayMs, int periodMs, int jitterMs) {
void schedule(Runnable action, int delayMs, int periodMs, int jitterMs) {
long initialDeadlineMs = time.milliseconds() + delayMs;
int eventId = eventIdGenerator.incrementAndGet();
PeriodicEvent event = new PeriodicEvent(action, eventId, random, initialDeadlineMs, periodMs, jitterMs);
Expand Down