Explorar el Código

bin/ircd: minor changes

ghassmo hace 4 años
padre
commit
8cdd47bf1b
Se han modificado 4 ficheros con 28 adiciones y 28 borrados
  1. 1 1
      bin/ircd/ircd_config.toml
  2. 20 17
      bin/ircd/src/main.rs
  3. 2 2
      bin/ircd/src/settings.rs
  4. 5 8
      src/raft/consensus.rs

+ 1 - 1
bin/ircd/ircd_config.toml

@@ -2,7 +2,7 @@
 #rpc_listen="127.0.0.1:11055"
 
 ## IRC listen URL
-#irc_listen="tcp://127.0.0.1:11066"
+#irc_listen="127.0.0.1:11066"
 
 ## Sets Datastore Path
 #datastore="~/.config/ircd"

+ 20 - 17
bin/ircd/src/main.rs

@@ -1,21 +1,21 @@
 use async_std::{
-    net::TcpStream,
+    net::{TcpListener, TcpStream},
     sync::{Arc, Mutex},
 };
+
 use std::net::SocketAddr;
 
 use async_channel::{Receiver, Sender};
 use async_executor::Executor;
 use easy_parallel::Parallel;
-use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt, StreamExt};
-use log::{debug, info, warn};
+use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
+use log::{debug, error, info, warn};
 use simplelog::{ColorChoice, TermLogger, TerminalMode};
 use smol::future;
 use structopt_toml::StructOptToml;
 
 use darkfi::{
     async_daemonize, net,
-    net::transport::{TcpTransport, Transport},
     raft::{NetMsg, ProtocolRaft, Raft},
     rpc::rpcserver::{listen_and_serve, RpcServerConfig},
     util::{
@@ -182,21 +182,26 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     //
     // IRC instance
     //
+    let listener = TcpListener::bind(settings.irc_listen).await?;
+    let local_addr = listener.local_addr()?;
+    info!("IRC listening on {}", local_addr);
     let executor_cloned = executor.clone();
+    let raft_receiver_cloned = raft_receiver.clone();
     let irc_task: smol::Task<Result<()>> = executor.spawn(async move {
-        let irc_listen_url = url::Url::parse(&settings.irc_listen)?;
-
-        let transport = TcpTransport::new(None, 1024);
-        let listener = transport.listen_on(irc_listen_url.clone()).unwrap().await.unwrap();
-        let mut incoming = listener.incoming();
-        info!("IRC start a TCP connection {}", &irc_listen_url.to_string());
-        while let Some(stream) = incoming.next().await {
-            let stream = stream.unwrap();
-            let peer_addr = stream.peer_addr()?;
-            info!("IRC Accepted TCP connection {}", peer_addr);
+        loop {
+            let (stream, peer_addr) = match listener.accept().await {
+                Ok((s, a)) => (s, a),
+                Err(e) => {
+                    error!("Failed listening for connections: {}", e);
+                    return Err(Error::ServiceStopped)
+                }
+            };
+
+            info!("IRC Accepted client: {}", peer_addr);
+
             executor_cloned
                 .spawn(process(
-                    raft_receiver.clone(),
+                    raft_receiver_cloned.clone(),
                     stream,
                     peer_addr,
                     raft_sender.clone(),
@@ -204,8 +209,6 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
                 ))
                 .detach();
         }
-
-        Ok(())
     });
 
     // Run once receive exit signal

+ 2 - 2
bin/ircd/src/settings.rs

@@ -21,8 +21,8 @@ pub struct Args {
     #[structopt(long = "rpc", default_value = "127.0.0.1:11055")]
     pub rpc_listen: SocketAddr,
     /// IRC listen URL
-    #[structopt(long = "irc", default_value = "tcp://127.0.0.1:11066")]
-    pub irc_listen: String,
+    #[structopt(long = "irc", default_value = "127.0.0.1:11066")]
+    pub irc_listen: SocketAddr,
     /// Sets Datastore Path
     #[structopt(long, default_value = "~/.config/ircd")]
     pub datastore: String,

+ 5 - 8
src/raft/consensus.rs

@@ -48,9 +48,9 @@ async fn load_node_ids_loop(
     }
 }
 
-async fn p2p_send_loop(p2p_recv: async_channel::Receiver<NetMsg>, p2p: net::P2pPtr) -> Result<()> {
+async fn p2p_send_loop(receiver: async_channel::Receiver<NetMsg>, p2p: net::P2pPtr) -> Result<()> {
     loop {
-        let msg: NetMsg = match p2p_recv.recv().await {
+        let msg: NetMsg = match receiver.recv().await {
             Ok(m) => m,
             Err(e) => {
                 error!(target: "raft", "error occurred while receiving a msg: {}", e);
@@ -346,13 +346,10 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
             self.set_current_term(&self.logs.0.last().unwrap().term.clone())?;
         }
 
-        if self.commit_length > sr.commit_length {
-            self.set_commit_length(&0)?;
-        }
-
         for i in self.commit_length..sr.commit_length {
             self.push_commit(&self.logs.get(i)?.msg).await?;
         }
+
         self.set_commit_length(&sr.commit_length)?;
 
         self.current_leader = Some(sr.leader_id.clone());
@@ -380,12 +377,12 @@ impl<T: Decodable + Encodable + Clone> Raft<T> {
 
     async fn waiting_for_sync(
         &mut self,
-        receive_queues: async_channel::Receiver<NetMsg>,
+        p2p_recv_channel: async_channel::Receiver<NetMsg>,
         stop_signal: async_channel::Receiver<()>,
     ) -> Result<()> {
         loop {
             select! {
-                msg =  receive_queues.recv().fuse() => {
+                msg =  p2p_recv_channel.recv().fuse() => {
                     let msg = msg?;
                     if msg.method == NetMsgMethod::SyncResponse {
                         let sr: SyncResponse = deserialize(&msg.payload)?;