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)

Reply via email to