Jelajahi Sumber

raft: add sent messages to seen_msgs vector

ghassmo 4 tahun lalu
induk
melakukan
9fdc11bd70
5 mengubah file dengan 29 tambahan dan 45 penghapusan
  1. 0 29
      Cargo.lock
  2. 2 2
      Cargo.toml
  3. 8 4
      bin/tau/taud/src/main.rs
  4. 15 6
      src/raft/consensus.rs
  5. 4 4
      src/raft/protocol_raft.rs

+ 0 - 29
Cargo.lock

@@ -2236,35 +2236,6 @@ dependencies = [
  "cfg-if 1.0.0",
  "cfg-if 1.0.0",
 ]
 ]
 
 
-[[package]]
-name = "irc-raft"
-version = "0.3.0"
-dependencies = [
- "async-channel",
- "async-executor",
- "async-std",
- "async-trait",
- "bs58",
- "clap 3.1.18",
- "crypto_box",
- "ctrlc-async",
- "darkfi",
- "easy-parallel",
- "futures",
- "futures-rustls",
- "fxhash",
- "log",
- "rand",
- "serde",
- "serde_json",
- "simplelog",
- "smol",
- "structopt",
- "structopt-toml",
- "toml",
- "url",
-]
-
 [[package]]
 [[package]]
 name = "ircd"
 name = "ircd"
 version = "0.3.0"
 version = "0.3.0"

+ 2 - 2
Cargo.toml

@@ -19,12 +19,12 @@ name = "darkfi"
 [workspace]
 [workspace]
 members = [
 members = [
 	"bin/zkas",
 	"bin/zkas",
-#"bin/cashierd",
+	#"bin/cashierd",
 	"bin/darkfid",
 	"bin/darkfid",
 	"bin/drk",
 	"bin/drk",
 	"bin/faucetd",
 	"bin/faucetd",
 	"bin/ircd",
 	"bin/ircd",
-	"bin/irc-raft",
+	#"bin/irc-raft",
 	"bin/dnetview",
 	"bin/dnetview",
 	"bin/daod",
 	"bin/daod",
 	"bin/dao-cli",
 	"bin/dao-cli",

+ 8 - 4
bin/tau/taud/src/main.rs

@@ -164,9 +164,14 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     //
     //
 
 
     let net_settings = settings.net;
     let net_settings = settings.net;
+    let seen_net_msgs = Arc::new(Mutex::new(vec![]));
 
 
     let datastore_raft = datastore_path.join("tau.db");
     let datastore_raft = datastore_path.join("tau.db");
-    let mut raft = Raft::<EncryptedTask>::new(net_settings.inbound.clone(), datastore_raft)?;
+    let mut raft = Raft::<EncryptedTask>::new(
+        net_settings.inbound.clone(),
+        datastore_raft,
+        seen_net_msgs.clone(),
+    )?;
 
 
     executor
     executor
         .spawn(start_sync_loop(
         .spawn(start_sync_loop(
@@ -187,15 +192,14 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
 
 
     let registry = p2p.protocol_registry();
     let registry = p2p.protocol_registry();
 
 
-    let seen_net_msg = Arc::new(Mutex::new(vec![]));
     let raft_node_id = raft.id.clone();
     let raft_node_id = raft.id.clone();
     registry
     registry
         .register(net::SESSION_ALL, move |channel, p2p| {
         .register(net::SESSION_ALL, move |channel, p2p| {
             let raft_node_id = raft_node_id.clone();
             let raft_node_id = raft_node_id.clone();
             let sender = p2p_send_channel.clone();
             let sender = p2p_send_channel.clone();
-            let seen_net_msg_cloned = seen_net_msg.clone();
+            let seen_net_msgs_cloned = seen_net_msgs.clone();
             async move {
             async move {
-                ProtocolRaft::init(raft_node_id, channel, sender, p2p, seen_net_msg_cloned).await
+                ProtocolRaft::init(raft_node_id, channel, sender, p2p, seen_net_msgs_cloned).await
             }
             }
         })
         })
         .await;
         .await;

+ 15 - 6
src/raft/consensus.rs

@@ -24,9 +24,9 @@ use super::{
     DataStore,
     DataStore,
 };
 };
 
 
-const HEARTBEATTIMEOUT: u64 = 1000;
-const TIMEOUT: u64 = 3000;
-const TIMEOUT_NODES: u64 = 3000;
+const HEARTBEATTIMEOUT: u64 = 500;
+const TIMEOUT: u64 = 6000;
+const TIMEOUT_NODES: u64 = 1000;
 
 
 async fn load_node_ids_loop(
 async fn load_node_ids_loop(
     nodes: Arc<Mutex<HashMap<NodeId, Url>>>,
     nodes: Arc<Mutex<HashMap<NodeId, Url>>>,
@@ -38,7 +38,7 @@ async fn load_node_ids_loop(
     }
     }
     loop {
     loop {
         debug!(target: "raft", "Loading node ids from p2p hosts",);
         debug!(target: "raft", "Loading node ids from p2p hosts",);
-        task::sleep(Duration::from_millis(TIMEOUT_NODES * 10)).await;
+        task::sleep(Duration::from_millis(TIMEOUT_NODES)).await;
         let hosts = p2p.hosts().clone();
         let hosts = p2p.hosts().clone();
         let nodes_ip = hosts.load_all().await.clone();
         let nodes_ip = hosts.load_all().await.clone();
 
 
@@ -87,10 +87,16 @@ pub struct Raft<T> {
     commits_channel: Channel<T>,
     commits_channel: Channel<T>,
 
 
     datastore: DataStore<T>,
     datastore: DataStore<T>,
+
+    seen_msgs: Arc<Mutex<Vec<u64>>>,
 }
 }
 
 
 impl<T: Decodable + Encodable + Clone> Raft<T> {
 impl<T: Decodable + Encodable + Clone> Raft<T> {
-    pub fn new(addr: Option<Url>, db_path: PathBuf) -> Result<Self> {
+    pub fn new(
+        addr: Option<Url>,
+        db_path: PathBuf,
+        seen_msgs: Arc<Mutex<Vec<u64>>>,
+    ) -> Result<Self> {
         if db_path.to_str().is_none() {
         if db_path.to_str().is_none() {
             error!(target: "raft", "datastore path is incorrect");
             error!(target: "raft", "datastore path is incorrect");
             return Err(Error::ParseFailed("unable to parse pathbuf to str"))
             return Err(Error::ParseFailed("unable to parse pathbuf to str"))
@@ -130,6 +136,7 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
             msgs_channel,
             msgs_channel,
             commits_channel,
             commits_channel,
             datastore,
             datastore,
+            seen_msgs,
         })
         })
     }
     }
 
 
@@ -348,6 +355,7 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
         self.role, random_id, &recipient_id.is_some(), &method);
         self.role, random_id, &recipient_id.is_some(), &method);
 
 
         let net_msg = NetMsg { id: random_id, recipient_id, payload: payload.to_vec(), method };
         let net_msg = NetMsg { id: random_id, recipient_id, payload: payload.to_vec(), method };
+        self.seen_msgs.lock().await.push(random_id);
         self.sender.0.send(net_msg).await?;
         self.sender.0.send(net_msg).await?;
 
 
         Ok(())
         Ok(())
@@ -395,6 +403,7 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
         self.set_current_term(&(self.current_term + 1))?;
         self.set_current_term(&(self.current_term + 1))?;
         self.role = Role::Candidate;
         self.role = Role::Candidate;
         self.set_voted_for(&Some(self_id.clone()))?;
         self.set_voted_for(&Some(self_id.clone()))?;
+        self.votes_received = vec![];
         self.votes_received.push(self_id.clone());
         self.votes_received.push(self_id.clone());
 
 
         self.reset_last_term();
         self.reset_last_term();
@@ -461,7 +470,7 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
             let nodes_cloned = nodes.clone();
             let nodes_cloned = nodes.clone();
             drop(nodes);
             drop(nodes);
 
 
-            if self.votes_received.len() >= ((nodes_cloned.len() + 1) / 2) {
+            if self.votes_received.len() >= (nodes_cloned.len() / 2) {
                 self.role = Role::Leader;
                 self.role = Role::Leader;
                 self.current_leader = Some(self.id.clone().unwrap());
                 self.current_leader = Some(self.id.clone().unwrap());
                 for node in nodes_cloned.iter() {
                 for node in nodes_cloned.iter() {

+ 4 - 4
src/raft/protocol_raft.rs

@@ -14,7 +14,7 @@ pub struct ProtocolRaft {
     notify_queue_sender: async_channel::Sender<NetMsg>,
     notify_queue_sender: async_channel::Sender<NetMsg>,
     msg_sub: net::MessageSubscription<NetMsg>,
     msg_sub: net::MessageSubscription<NetMsg>,
     p2p: net::P2pPtr,
     p2p: net::P2pPtr,
-    msgs: Arc<Mutex<Vec<u64>>>,
+    seen_msgs: Arc<Mutex<Vec<u64>>>,
 }
 }
 
 
 impl ProtocolRaft {
 impl ProtocolRaft {
@@ -23,7 +23,7 @@ impl ProtocolRaft {
         channel: net::ChannelPtr,
         channel: net::ChannelPtr,
         notify_queue_sender: async_channel::Sender<NetMsg>,
         notify_queue_sender: async_channel::Sender<NetMsg>,
         p2p: net::P2pPtr,
         p2p: net::P2pPtr,
-        msgs: Arc<Mutex<Vec<u64>>>,
+        seen_msgs: Arc<Mutex<Vec<u64>>>,
     ) -> net::ProtocolBasePtr {
     ) -> net::ProtocolBasePtr {
         let message_subsytem = channel.get_message_subsystem();
         let message_subsytem = channel.get_message_subsystem();
         message_subsytem.add_dispatch::<NetMsg>().await;
         message_subsytem.add_dispatch::<NetMsg>().await;
@@ -36,7 +36,7 @@ impl ProtocolRaft {
             msg_sub,
             msg_sub,
             jobsman: net::ProtocolJobsManager::new("ProtocolRaft", channel),
             jobsman: net::ProtocolJobsManager::new("ProtocolRaft", channel),
             p2p,
             p2p,
-            msgs,
+            seen_msgs,
         })
         })
     }
     }
 
 
@@ -52,7 +52,7 @@ impl ProtocolRaft {
             );
             );
 
 
             {
             {
-                let mut msgs = self.msgs.lock().await;
+                let mut msgs = self.seen_msgs.lock().await;
                 if msgs.contains(&msg.id) {
                 if msgs.contains(&msg.id) {
                     continue
                     continue
                 }
                 }