Просмотр исходного кода

bin/ircd: redesign on main.rs code for handling connections

ghassmo 3 лет назад
Родитель
Сommit
bd9789850e
4 измененных файлов с 130 добавлено и 123 удалено
  1. 0 0
      bin/ircd/src/irc_server/command.rs
  2. 1 1
      bin/ircd/src/irc_server/mod.rs
  3. 127 122
      bin/ircd/src/main.rs
  4. 2 0
      bin/ircd/src/settings.rs

+ 0 - 0
bin/ircd/src/server/command.rs → bin/ircd/src/irc_server/command.rs


+ 1 - 1
bin/ircd/src/server/mod.rs → bin/ircd/src/irc_server/mod.rs

@@ -18,7 +18,7 @@ mod command;
 pub struct IrcServerConnection<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> {
     // server stream
     write_stream: WriteHalf<C>,
-    peer_address: SocketAddr,
+    pub peer_address: SocketAddr,
     // msg ids
     seen_msg_ids: SeenMsgIds,
     privmsgs_buffer: ArcPrivmsgsBuffer,

+ 127 - 122
bin/ircd/src/main.rs

@@ -7,7 +7,10 @@ use std::{fmt, fs::File, net::SocketAddr};
 
 use async_channel::Receiver;
 use async_executor::Executor;
-use futures::{io::BufReader, AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWrite, FutureExt};
+use futures::{
+    io::{BufReader, ReadHalf},
+    AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWrite, FutureExt,
+};
 use futures_rustls::{rustls, TlsAcceptor};
 use fxhash::FxHashMap;
 use log::{error, info, warn};
@@ -18,7 +21,7 @@ use structopt_toml::StructOptToml;
 use darkfi::{
     async_daemonize, net,
     rpc::server::listen_and_serve,
-    system::{Subscriber, SubscriberPtr},
+    system::{Subscriber, SubscriberPtr, Subscription},
     util::{
         cli::{get_log_config, get_log_level, spawn_config},
         expand_path,
@@ -30,18 +33,18 @@ use darkfi::{
 
 pub mod buffers;
 pub mod crypto;
+pub mod irc_server;
 pub mod privmsg;
 pub mod protocol_privmsg;
 pub mod rpc;
-pub mod server;
 pub mod settings;
 
 use crate::{
     buffers::{ArcPrivmsgsBuffer, PrivmsgsBuffer, RingBuffer, SeenMsgIds, SIZE_OF_MSG_IDSS_BUFFER},
+    irc_server::IrcServerConnection,
     privmsg::Privmsg,
     protocol_privmsg::ProtocolPrivmsg,
     rpc::JsonRpcInterface,
-    server::IrcServerConnection,
     settings::{
         parse_configured_channels, parse_configured_contacts, Args, ChannelInfo, CONFIG_FILE,
         CONFIG_FILE_CONTENTS,
@@ -63,6 +66,79 @@ impl fmt::Display for KeyPair {
     }
 }
 
+async fn setup_listener(settings: Args) -> Result<(TcpListener, Option<TlsAcceptor>)> {
+
+    let listenaddr = settings.irc_listen.socket_addrs(|| None)?[0];
+    let listener = TcpListener::bind(listenaddr).await?;
+
+    let acceptor = match settings.irc_listen.scheme() {
+        "tls" => {
+            // openssl genpkey -algorithm ED25519 > example.com.key
+            // openssl req -new -out example.com.csr -key example.com.key
+            // openssl x509 -req -days 700 -in example.com.csr -signkey example.com.key -out example.com.crt
+
+            if settings.irc_tls_secret.is_none() || settings.irc_tls_cert.is_none() {
+                error!("To listen using TLS, please set irc_tls_secret and irc_tls_cert in your config file.");
+                return Err(Error::KeypairPathNotFound)
+            }
+
+            let file = File::open(expand_path(&settings.irc_tls_secret.unwrap())?)?;
+            let mut reader = std::io::BufReader::new(file);
+            let secret = &rustls_pemfile::pkcs8_private_keys(&mut reader)?[0];
+            let secret = rustls::PrivateKey(secret.clone());
+
+            let file = File::open(expand_path(&settings.irc_tls_cert.unwrap())?)?;
+            let mut reader = std::io::BufReader::new(file);
+            let certificate = &rustls_pemfile::certs(&mut reader)?[0];
+            let certificate = rustls::Certificate(certificate.clone());
+
+            let config = rustls::ServerConfig::builder()
+                .with_safe_defaults()
+                .with_no_client_auth()
+                .with_single_cert(vec![certificate], secret)?;
+
+            let acceptor = TlsAcceptor::from(Arc::new(config));
+            Some(acceptor)
+        }
+        _ => None,
+    };
+    Ok((listener, acceptor))
+}
+
+async fn start_listening(ircd: Ircd, executor: Arc<Executor<'_>>, settings: Args) -> Result<()> {
+    let (listener, acceptor) = setup_listener(settings.clone()).await?;
+    info!("IRC listening on {}", settings.irc_listen);
+    loop {
+        let (stream, peer_addr) = match listener.accept().await {
+            Ok((s, a)) => (s, a),
+            Err(e) => {
+                error!("failed accepting new connections: {}", e);
+                continue
+            }
+        };
+
+        let result = if let Some(acceptor) = acceptor.clone() {
+            let stream = match acceptor.accept(stream).await {
+                Ok(s) => s,
+                Err(e) => {
+                    error!("Failed accepting TLS connection: {}", e);
+                    continue
+                }
+            };
+            ircd.process_new_connection(executor.clone(), stream, peer_addr).await
+        } else {
+            ircd.process_new_connection(executor.clone(), stream, peer_addr).await
+        };
+
+        if let Err(e) = result {
+            error!("Failed processing connection {}: {}", peer_addr, e);
+            continue
+        };
+
+        info!("IRC Accepted new client: {}", peer_addr);
+    }
+}
+
 struct Ircd {
     // msgs
     seen_msg_ids: SeenMsgIds,
@@ -119,13 +195,13 @@ impl Ircd {
     ) -> Result<()> {
         let (reader, writer) = stream.split();
 
-        let mut reader = BufReader::new(reader);
+        let reader = BufReader::new(reader);
 
         // New subscriber
         let receiver = self.senders.clone().subscribe().await;
 
         // New irc connection
-        let mut conn = IrcServerConnection::new(
+        let conn = IrcServerConnection::new(
             writer,
             peer_addr,
             self.seen_msg_ids.clone(),
@@ -139,40 +215,37 @@ impl Ircd {
             receiver.get_id(),
         );
 
-        executor
-            .spawn(async move {
-                loop {
-                    let mut line = String::new();
-
-                    let result: Result<()> = futures::select! {
-                        msg = receiver.receive().fuse() => {
-                            match conn.process_msg_from_p2p(&msg).await {
-                                Ok(_) => Ok(()),
-                                Err(e) => {
-                                    error!("Process msg from p2p failed {}: {}", peer_addr, e);
-                                    Err(Error::ChannelStopped)
-                                }
-                            }
-                        }
-                        err = reader.read_line(&mut line).fuse() => {
-                            match conn.process_line_from_client(err, line).await {
-                                Ok(_) => Ok(()),
-                                Err(e) => {
-                                    error!("Process line from client failed {}: {}", peer_addr, e);
-                                    Err(Error::ChannelStopped)
-                                }
-                            }
-                        }
-                    };
-
-                    if let Err(e) = result {
-                        warn!("Close connection for clinet {}: {}", peer_addr, e);
-                        receiver.unsubscribe().await;
+        executor.spawn(Self::listen(conn, reader, receiver)).detach();
+
+        Ok(())
+    }
+
+    async fn listen<C: AsyncRead + AsyncWrite + Send + Unpin + 'static>(
+        mut conn: IrcServerConnection<C>,
+        mut reader: BufReader<ReadHalf<C>>,
+        receiver: Subscription<Privmsg>,
+    ) -> Result<()> {
+        loop {
+            let mut line = String::new();
+
+            futures::select! {
+                msg = receiver.receive().fuse() => {
+                    if let Err(e) = conn.process_msg_from_p2p(&msg).await {
+                        error!("Process msg from p2p failed {}: {}",  conn.peer_address, e);
                         break
                     }
                 }
-            })
-            .detach();
+                err = reader.read_line(&mut line).fuse() => {
+                    if let Err(e) = conn.process_line_from_client(err, line).await {
+                        error!("Process line from client failed {}: {}", conn.peer_address, e);
+                        break
+                    }
+                }
+            }
+        }
+
+        warn!("Close connection for clinet {}", conn.peer_address);
+        receiver.unsubscribe().await;
 
         Ok(())
     }
@@ -208,10 +281,10 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
         return Ok(())
     }
 
-    let password = settings.password.unwrap_or_default();
+    let password = settings.password.clone().unwrap_or_default();
 
     // Pick up channel settings from the TOML configuration
-    let cfg_path = get_config_path(settings.config, CONFIG_FILE)?;
+    let cfg_path = get_config_path(settings.config.clone(), CONFIG_FILE)?;
     let toml_contents = std::fs::read_to_string(cfg_path)?;
     let configured_chans = parse_configured_channels(&toml_contents)?;
     let configured_contacts = parse_configured_contacts(&toml_contents)?;
@@ -219,7 +292,7 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     //
     // P2p setup
     //
-    let net_settings = settings.net;
+    let net_settings = settings.net.clone();
     let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<Privmsg>();
 
     let p2p = net::P2p::new(net_settings.into()).await;
@@ -252,6 +325,8 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     let executor_cloned = executor.clone();
     executor_cloned.spawn(p2p.clone().run(executor.clone())).detach();
 
+    p2p.clone().wait_for_outbound(executor.clone()).await?;
+
     //
     // RPC interface
     //
@@ -263,89 +338,19 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     //
     // IRC instance
     //
-    let listenaddr = settings.irc_listen.socket_addrs(|| None)?[0];
-    let listener = TcpListener::bind(listenaddr).await?;
-
-    let acceptor = match settings.irc_listen.scheme() {
-        "tls" => {
-            // openssl genpkey -algorithm ED25519 > example.com.key
-            // openssl req -new -out example.com.csr -key example.com.key
-            // openssl x509 -req -days 700 -in example.com.csr -signkey example.com.key -out example.com.crt
-
-            if settings.irc_tls_secret.is_none() || settings.irc_tls_cert.is_none() {
-                error!("To listen using TLS, please set irc_tls_secret and irc_tls_cert in your config file.");
-                return Err(Error::KeypairPathNotFound)
-            }
-
-            let file = File::open(expand_path(&settings.irc_tls_secret.unwrap())?)?;
-            let mut reader = std::io::BufReader::new(file);
-            let secret = &rustls_pemfile::pkcs8_private_keys(&mut reader)?[0];
-            let secret = rustls::PrivateKey(secret.clone());
-
-            let file = File::open(expand_path(&settings.irc_tls_cert.unwrap())?)?;
-            let mut reader = std::io::BufReader::new(file);
-            let certificate = &rustls_pemfile::certs(&mut reader)?[0];
-            let certificate = rustls::Certificate(certificate.clone());
 
-            let config = rustls::ServerConfig::builder()
-                .with_safe_defaults()
-                .with_no_client_auth()
-                .with_single_cert(vec![certificate], secret)?;
-
-            let acceptor = TlsAcceptor::from(Arc::new(config));
-            Some(acceptor)
-        }
-        _ => None,
-    };
-
-    info!("IRC listening on {}", settings.irc_listen);
-
-    let executor_cloned = executor.clone();
-    executor
-        .spawn(async move {
-            let ircd = Ircd::new(
-                seen_msg_ids.clone(),
-                privmsgs_buffer.clone(),
-                settings.autojoin.clone(),
-                password.clone(),
-                configured_chans.clone(),
-                configured_contacts.clone(),
-                p2p.clone(),
-            );
-
-            ircd.start_p2p_receive_loop(executor_cloned.clone(), p2p_recv_channel);
-
-            loop {
-                let (stream, peer_addr) = match listener.accept().await {
-                    Ok((s, a)) => (s, a),
-                    Err(e) => {
-                        error!("failed accepting new connections: {}", e);
-                        continue
-                    }
-                };
-
-                let result = if let Some(acceptor) = acceptor.clone() {
-                    let stream = match acceptor.accept(stream).await {
-                        Ok(s) => s,
-                        Err(e) => {
-                            error!("Failed accepting TLS connection: {}", e);
-                            continue
-                        }
-                    };
-                    ircd.process_new_connection(executor_cloned.clone(), stream, peer_addr).await
-                } else {
-                    ircd.process_new_connection(executor_cloned.clone(), stream, peer_addr).await
-                };
-
-                if let Err(e) = result {
-                    error!("Failed processing connection {}: {}", peer_addr, e);
-                    continue
-                };
-
-                info!("IRC Accepted new client: {}", peer_addr);
-            }
-        })
-        .detach();
+    let ircd = Ircd::new(
+        seen_msg_ids.clone(),
+        privmsgs_buffer.clone(),
+        settings.autojoin.clone(),
+        password.clone(),
+        configured_chans.clone(),
+        configured_contacts.clone(),
+        p2p.clone(),
+    );
+
+    ircd.start_p2p_receive_loop(executor.clone(), p2p_recv_channel);
+    executor.spawn(start_listening(ircd, executor.clone(), settings.clone())).detach();
 
     // Run once receive exit signal
     let (signal, shutdown) = async_channel::bounded::<()>(1);

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

@@ -126,6 +126,7 @@ fn parse_priv_key(data: &str) -> Result<String> {
     Ok(pk)
 }
 
+
 /// Parse a TOML string for any configured contact list and return
 /// a map containing said configurations.
 ///
@@ -166,6 +167,7 @@ pub fn parse_configured_contacts(data: &str) -> Result<FxHashMap<String, Contact
     Ok(ret)
 }
 
+/// TODO CLEAN UP
 /// Parse a TOML string for any configured channels and return
 /// a map containing said configurations.
 ///