|
@@ -204,11 +204,11 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
|
|
|
|
|
|
|
|
select! {
|
|
select! {
|
|
|
m = p2p_recv_channel.recv().fuse() => result = self.handle_method(m?).await,
|
|
m = p2p_recv_channel.recv().fuse() => result = self.handle_method(m?).await,
|
|
|
- m = broadcast_msg_rv.recv().fuse() => result = self.broadcast_msg(&m?,None).await,
|
|
|
|
|
|
|
+ m = broadcast_msg_rv.recv().fuse() => result = self.broadcast_msg(&m?, None).await,
|
|
|
_ = task::sleep(timeout).fuse() => {
|
|
_ = task::sleep(timeout).fuse() => {
|
|
|
result = if self.role == Role::Leader {
|
|
result = if self.role == Role::Leader {
|
|
|
self.send_heartbeat().await
|
|
self.send_heartbeat().await
|
|
|
- }else {
|
|
|
|
|
|
|
+ } else {
|
|
|
self.send_vote_request().await
|
|
self.send_vote_request().await
|
|
|
};
|
|
};
|
|
|
},
|
|
},
|
|
@@ -408,14 +408,13 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
async fn send_heartbeat(&self) -> Result<()> {
|
|
async fn send_heartbeat(&self) -> Result<()> {
|
|
|
- if self.role == Role::Leader {
|
|
|
|
|
- let nodes = self.nodes.lock().await;
|
|
|
|
|
- let nodes_cloned = nodes.clone();
|
|
|
|
|
- drop(nodes);
|
|
|
|
|
- for node in nodes_cloned.iter() {
|
|
|
|
|
- self.update_logs(node.0).await?;
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ let nodes = self.nodes.lock().await;
|
|
|
|
|
+ let nodes_cloned = nodes.clone();
|
|
|
|
|
+ drop(nodes);
|
|
|
|
|
+ for node in nodes_cloned.iter() {
|
|
|
|
|
+ self.update_logs(node.0).await?;
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
Ok(())
|
|
Ok(())
|
|
|
}
|
|
}
|
|
|
|
|
|