Jelajahi Sumber

raft: avoid broadcasting when the role is Candidate

ghassmo 4 tahun lalu
induk
melakukan
afdd4b45dc
2 mengubah file dengan 28 tambahan dan 15 penghapusan
  1. 24 15
      src/raft/consensus.rs
  2. 4 0
      src/raft/consensus_candidate.rs

+ 24 - 15
src/raft/consensus.rs

@@ -208,21 +208,30 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
     }
 
     async fn broadcast_msg(&mut self, msg: &T, msg_id: Option<u64>) -> Result<()> {
-        if self.role == Role::Leader {
-            let msg = serialize(msg);
-            let log = Log { msg, term: self.current_term()? };
-            self.push_log(&log)?;
-
-            self.acked_length.insert(&self.id, self.logs_len());
-        } else {
-            let b_msg = BroadcastMsgRequest(serialize(msg));
-            self.send(
-                Some(self.current_leader.clone()),
-                &serialize(&b_msg),
-                NetMsgMethod::BroadcastRequest,
-                msg_id,
-            )
-            .await?;
+        loop {
+            match self.role {
+                Role::Leader => {
+                    let msg = serialize(msg);
+                    let log = Log { msg, term: self.current_term()? };
+                    self.push_log(&log)?;
+                    self.acked_length.insert(&self.id, self.logs_len());
+                    break
+                }
+                Role::Follower => {
+                    let b_msg = BroadcastMsgRequest(serialize(msg));
+                    self.send(
+                        Some(self.current_leader.clone()),
+                        &serialize(&b_msg),
+                        NetMsgMethod::BroadcastRequest,
+                        msg_id,
+                    )
+                    .await?;
+                    break
+                }
+                Role::Candidate => {
+                    util::sleep(2).await;
+                }
+            }
         }
 
         debug!(target: "raft", "Role: {:?} Id: {:?}, broadcast a msg id: {:?} ", self.role, self.id, msg_id);

+ 4 - 0
src/raft/consensus_candidate.rs

@@ -1,3 +1,5 @@
+use log::info;
+
 use crate::{
     util::serial::{serialize, Decodable, Encodable},
     Result,
@@ -13,6 +15,7 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
         let self_id = self.id();
 
         self.set_current_term(&(self.current_term()? + 1))?;
+        info!(target: "raft", "Set the node role as Candidate");
         self.role = Role::Candidate;
         self.set_voted_for(&Some(self_id.clone()))?;
         self.votes_received = vec![];
@@ -40,6 +43,7 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
             drop(nodes);
 
             if self.votes_received.len() >= ((nodes_cloned.len() + 1) / 2) {
+                info!(target: "raft", "Set the node role as Leader");
                 self.role = Role::Leader;
                 self.current_leader = self.id();
                 for node in nodes_cloned.iter() {