This is an automated email from the ASF dual-hosted git repository.
tanxinyu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new a0031aece9 Consensus CI test fix transferLeader assert (#6264)
a0031aece9 is described below
commit a0031aece9069ba30ce400d1c640d1251b0ff929
Author: William Song <[email protected]>
AuthorDate: Fri Jun 17 11:18:33 2022 +0800
Consensus CI test fix transferLeader assert (#6264)
* consensus CI test fix transferLeader assert
* add comments
* rerun CI
* test CI params
* test CI params
* refactor ratis consensus unit test
---
.../iotdb/consensus/ratis/RatisConsensus.java | 1 +
.../iotdb/consensus/ratis/RatisConsensusTest.java | 120 ++++++++++-----------
2 files changed, 58 insertions(+), 63 deletions(-)
diff --git
a/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
b/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
index 726012a3ed..8687dba132 100644
---
a/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
+++
b/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
@@ -410,6 +410,7 @@ class RatisConsensus implements IConsensus {
}
/**
+ * NOTICE: transferLeader *does not guarantee* the leader be transferred to
newLeader.
* transferLeader is implemented by 1. modify peer priority 2. ask current
leader to step down
*
* <p>1. call setConfiguration to upgrade newLeader's priority to 1 and
degrade all follower peers
diff --git
a/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
b/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
index 063c798722..456a27fdde 100644
---
a/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
+++
b/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
@@ -29,6 +29,7 @@ import
org.apache.iotdb.consensus.common.request.ByteBufferConsensusRequest;
import org.apache.iotdb.consensus.common.response.ConsensusReadResponse;
import org.apache.iotdb.consensus.common.response.ConsensusWriteResponse;
import org.apache.iotdb.consensus.config.ConsensusConfig;
+import org.apache.iotdb.consensus.config.RatisConfig;
import org.apache.ratis.util.FileUtils;
import org.junit.After;
@@ -40,7 +41,6 @@ import java.io.File;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
-import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
@@ -60,12 +60,22 @@ public class RatisConsensusTest {
private void makeServers() throws IOException {
for (int i = 0; i < 3; i++) {
stateMachines.add(new TestUtils.IntegerCounter());
+ RatisConfig config =
+ RatisConfig.newBuilder()
+ .setLog(
+ RatisConfig.Log.newBuilder()
+ .setPurgeUptoSnapshotIndex(true)
+ .setPurgeGap(10)
+ .build())
+
.setSnapshot(RatisConfig.Snapshot.newBuilder().setAutoTriggerThreshold(100).build())
+ .build();
int finalI = i;
servers.add(
ConsensusFactory.getConsensusImpl(
ConsensusFactory.RatisConsensus,
ConsensusConfig.newBuilder()
.setThisNode(peers.get(i).getEndpoint())
+ .setRatisConfig(config)
.setStorageDir(peersStorage.get(i).getAbsolutePath())
.build(),
groupId -> stateMachines.get(finalI))
@@ -114,84 +124,68 @@ public class RatisConsensusTest {
}
@Test
- public void basicConsensus() throws Exception {
+ public void basicConsensus3Copy() throws Exception {
+ servers.get(0).addConsensusGroup(group.getGroupId(), group.getPeers());
+ servers.get(1).addConsensusGroup(group.getGroupId(), group.getPeers());
+ servers.get(2).addConsensusGroup(group.getGroupId(), group.getPeers());
+
+ doConsensus(servers.get(0), group.getGroupId(), 10, 10);
+ }
+
+ @Test
+ public void addMemberToGroup() throws Exception {
+ List<Peer> original = peers.subList(0, 1);
+
+ servers.get(0).addConsensusGroup(group.getGroupId(), original);
+ doConsensus(servers.get(0), group.getGroupId(), 10, 10);
+
+ // add 2 members
+ servers.get(1).addConsensusGroup(group.getGroupId(), peers);
+ servers.get(0).addPeer(group.getGroupId(), peers.get(1));
+
+ servers.get(2).addConsensusGroup(group.getGroupId(), peers);
+ servers.get(0).changePeer(group.getGroupId(), peers);
- // 1. Add a new group
+ Assert.assertEquals(stateMachines.get(0).getConfiguration().size(), 3);
+ doConsensus(servers.get(0), group.getGroupId(), 10, 20);
+ }
+
+ @Test
+ public void removeMemberFromGroup() throws Exception {
servers.get(0).addConsensusGroup(group.getGroupId(), group.getPeers());
servers.get(1).addConsensusGroup(group.getGroupId(), group.getPeers());
servers.get(2).addConsensusGroup(group.getGroupId(), group.getPeers());
- // 2. Do Consensus 10
doConsensus(servers.get(0), group.getGroupId(), 10, 10);
- // pick a random leader
- int leader = 1;
- int follower1 = (leader + 1) % 3;
- int follower2 = (leader + 2) % 3;
-
- // 3. Remove two Peers from Group (peer 0 and peer 2)
- // transfer the leader to peer1
- servers.get(follower1).transferLeader(gid, peers.get(leader));
- // Assert.assertTrue(servers.get(1).isLeader(gid));
- // first use removePeer to inform the group leader of configuration change
- servers.get(leader).removePeer(gid, peers.get(follower1));
- servers.get(leader).removePeer(gid, peers.get(follower2));
- // then use removeConsensusGroup to clean up removed Consensus-Peer's
states
- servers.get(follower1).removeConsensusGroup(gid);
- servers.get(follower2).removeConsensusGroup(gid);
- Assert.assertEquals(
- servers.get(leader).getLeader(gid).getEndpoint(),
peers.get(leader).getEndpoint());
- Assert.assertEquals(
- stateMachines.get(leader).getLeaderEndpoint(),
peers.get(leader).getEndpoint());
- Assert.assertEquals(stateMachines.get(leader).getConfiguration().size(),
1);
- Assert.assertEquals(stateMachines.get(leader).getConfiguration().get(0),
peers.get(leader));
-
- // 4. try consensus again with one peer
- doConsensus(servers.get(leader), gid, 10, 20);
-
- // 5. add two peers back
- // first notify these new peers, let them initialize
- servers.get(follower1).addConsensusGroup(gid, peers);
- servers.get(follower2).addConsensusGroup(gid, peers);
- // then use addPeer to inform the group leader of configuration change
- servers.get(leader).addPeer(gid, peers.get(follower1));
- servers.get(leader).addPeer(gid, peers.get(follower2));
- Assert.assertEquals(stateMachines.get(leader).getConfiguration().size(),
3);
-
- // 6. try consensus with all 3 peers
- doConsensus(servers.get(2), gid, 10, 30);
-
- // pick a random leader
- leader = 0;
- follower1 = (leader + 1) % 3;
- follower2 = (leader + 2) % 3;
- // 7. again, group contains only peer0
- servers.get(0).transferLeader(group.getGroupId(), peers.get(leader));
- servers
- .get(leader)
- .changePeer(group.getGroupId(),
Collections.singletonList(peers.get(leader)));
- servers.get(follower1).removeConsensusGroup(group.getGroupId());
- servers.get(follower2).removeConsensusGroup(group.getGroupId());
- Assert.assertEquals(
- stateMachines.get(leader).getLeaderEndpoint(),
peers.get(leader).getEndpoint());
- Assert.assertEquals(stateMachines.get(leader).getConfiguration().size(),
1);
- Assert.assertEquals(stateMachines.get(leader).getConfiguration().get(0),
peers.get(leader));
-
- // 8. try consensus with only peer0
- doConsensus(servers.get(leader), gid, 10, 40);
-
- // 9. shutdown all the servers
+ servers.get(0).transferLeader(gid, peers.get(0));
+ servers.get(0).removePeer(gid, peers.get(1));
+ servers.get(1).removeConsensusGroup(gid);
+ servers.get(0).removePeer(gid, peers.get(2));
+ servers.get(2).removeConsensusGroup(gid);
+
+ doConsensus(servers.get(0), group.getGroupId(), 10, 20);
+ }
+
+ @Test
+ public void crashAndStart() throws Exception {
+ servers.get(0).addConsensusGroup(group.getGroupId(), group.getPeers());
+ servers.get(1).addConsensusGroup(group.getGroupId(), group.getPeers());
+ servers.get(2).addConsensusGroup(group.getGroupId(), group.getPeers());
+
+ // 200 operation will trigger snapshot & purge
+ doConsensus(servers.get(0), group.getGroupId(), 200, 200);
+
for (IConsensus consensus : servers) {
consensus.stop();
}
servers.clear();
- // 10. start again and verify the snapshot
makeServers();
servers.get(0).addConsensusGroup(group.getGroupId(), group.getPeers());
servers.get(1).addConsensusGroup(group.getGroupId(), group.getPeers());
servers.get(2).addConsensusGroup(group.getGroupId(), group.getPeers());
- doConsensus(servers.get(0), gid, 10, 50);
+ doConsensus(servers.get(0), gid, 10, 210);
}
private void doConsensus(IConsensus consensus, ConsensusGroupId gid, int
count, int target)