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

bin/ircd: combine all irc config in one struct & create Buffers struct for read/unread msgs

ghassmo 3 лет назад
Родитель
Сommit
85834f9e27
6 измененных файлов с 244 добавлено и 279 удалено
  1. 47 7
      bin/ircd/src/buffers.rs
  2. 11 17
      bin/ircd/src/crypto.rs
  3. 84 111
      bin/ircd/src/irc/client.rs
  4. 69 33
      bin/ircd/src/irc/mod.rs
  5. 17 83
      bin/ircd/src/main.rs
  6. 16 28
      bin/ircd/src/protocol_privmsg.rs

+ 47 - 7
bin/ircd/src/buffers.rs

@@ -2,6 +2,8 @@ use async_std::sync::{Arc, Mutex};
 use std::{cmp::Ordering, collections::VecDeque};
 
 use chrono::Utc;
+use fxhash::FxHashMap;
+use ripemd::{Digest, Ripemd160};
 
 use crate::Privmsg;
 
@@ -9,6 +11,48 @@ pub const SIZE_OF_MSGS_BUFFER: usize = 4095;
 pub const SIZE_OF_MSG_IDSS_BUFFER: usize = 65536;
 pub const LIFETIME_FOR_ORPHAN: i64 = 600;
 
+pub type InvSeenIds = Arc<Mutex<RingBuffer<u64>>>;
+pub type SeenIds = Mutex<RingBuffer<u64>>;
+pub type MutexPrivmsgsBuffer = Mutex<PrivmsgsBuffer>;
+pub type UnreadMsgs = Mutex<UMsgs>;
+pub type Buffers = Arc<Msgs>;
+
+pub struct Msgs {
+    pub privmsgs: MutexPrivmsgsBuffer,
+    pub unread_msgs: UnreadMsgs,
+    pub seen_ids: SeenIds,
+}
+
+pub fn create_buffers() -> Buffers {
+    let seen_ids = Mutex::new(RingBuffer::new(SIZE_OF_MSG_IDSS_BUFFER));
+    let privmsgs = PrivmsgsBuffer::new();
+    let unread_msgs = Mutex::new(UMsgs::new());
+    Arc::new(Msgs { privmsgs, unread_msgs, seen_ids })
+}
+
+#[derive(Clone)]
+pub struct UMsgs(pub FxHashMap<String, Privmsg>);
+
+impl UMsgs {
+    pub fn new() -> Self {
+        Self(FxHashMap::default())
+    }
+
+    pub fn insert(&mut self, msg: &Privmsg) -> String {
+        let mut hasher = Ripemd160::new();
+        hasher.update(msg.to_string());
+        let key = hex::encode(hasher.finalize());
+        self.0.insert(key.clone(), msg.clone());
+        key
+    }
+}
+
+impl Default for UMsgs {
+    fn default() -> Self {
+        Self::new()
+    }
+}
+
 #[derive(Clone)]
 pub struct RingBuffer<T> {
     pub items: VecDeque<T>,
@@ -56,21 +100,17 @@ impl<T: Eq + PartialEq + Clone> RingBuffer<T> {
     }
 }
 
-pub type SeenIds = Arc<Mutex<RingBuffer<u64>>>;
-
-pub type ArcPrivmsgsBuffer = Arc<Mutex<PrivmsgsBuffer>>;
-
 pub struct PrivmsgsBuffer {
     buffer: RingBuffer<Privmsg>,
     orphans: RingBuffer<Orphan>,
 }
 
 impl PrivmsgsBuffer {
-    pub fn new() -> ArcPrivmsgsBuffer {
-        Arc::new(Mutex::new(Self {
+    pub fn new() -> MutexPrivmsgsBuffer {
+        Mutex::new(Self {
             buffer: RingBuffer::new(SIZE_OF_MSGS_BUFFER),
             orphans: RingBuffer::new(SIZE_OF_MSGS_BUFFER),
-        }))
+        })
     }
 
     pub fn push(&mut self, privmsg: &Privmsg) {

+ 11 - 17
bin/ircd/src/crypto.rs

@@ -6,13 +6,12 @@ use fxhash::FxHashMap;
 use rand::rngs::OsRng;
 
 use crate::{
-    privmsg::{Privmsg, MAXIMUM_LENGTH_OF_NICKNAME},
     settings::{ChannelInfo, ContactInfo},
+    Privmsg,
 };
 
-/// Try decrypting a message given a NaCl box and a base58 string.
 /// The format we're using is nonce+ciphertext, where nonce is 24 bytes.
-fn try_decrypt_message(salt_box: &SalsaBox, ciphertext: &str) -> Option<String> {
+fn try_decrypt(salt_box: &SalsaBox, ciphertext: &str) -> Option<String> {
     let bytes = match bs58::decode(ciphertext).into_vec() {
         Ok(v) => v,
         Err(_) => return None,
@@ -38,9 +37,8 @@ fn try_decrypt_message(salt_box: &SalsaBox, ciphertext: &str) -> Option<String>
     }
 }
 
-/// Encrypt a message given a NaCl box and a plaintext string.
 /// The format we're using is nonce+ciphertext, where nonce is 24 bytes.
-pub fn encrypt_message(salt_box: &SalsaBox, plaintext: &str) -> String {
+pub fn encrypt(salt_box: &SalsaBox, plaintext: &str) -> String {
     let nonce = SalsaBox::generate_nonce(&mut OsRng);
     let mut ciphertext = salt_box.encrypt(&nonce, plaintext.as_bytes()).unwrap();
 
@@ -67,7 +65,7 @@ pub fn decrypt_target(
         let salt_box = chan_info.salt_box.clone();
 
         if let Some(salt_box) = salt_box {
-            let decrypted_target = try_decrypt_message(&salt_box, &privmsg.target);
+            let decrypted_target = try_decrypt(&salt_box, &privmsg.target);
             if decrypted_target.is_none() {
                 continue
             }
@@ -85,7 +83,7 @@ pub fn decrypt_target(
 
         let salt_box = cnt_info.salt_box.clone();
         if let Some(salt_box) = salt_box {
-            let decrypted_target = try_decrypt_message(&salt_box, &privmsg.target);
+            let decrypted_target = try_decrypt(&salt_box, &privmsg.target);
             if decrypted_target.is_none() {
                 continue
             }
@@ -100,23 +98,19 @@ pub fn decrypt_target(
 
 /// Decrypt PrivMsg nickname and message
 pub fn decrypt_privmsg(salt_box: &SalsaBox, privmsg: &mut Privmsg) {
-    let decrypted_nick = try_decrypt_message(&salt_box.clone(), &privmsg.nickname);
-    let decrypted_msg = try_decrypt_message(&salt_box.clone(), &privmsg.message);
+    let decrypted_nick = try_decrypt(&salt_box.clone(), &privmsg.nickname);
+    let decrypted_msg = try_decrypt(&salt_box.clone(), &privmsg.message);
 
-    if decrypted_nick.is_none() | decrypted_msg.is_none() {
+    if decrypted_nick.is_none() && decrypted_msg.is_none() {
         return
     }
-
     privmsg.nickname = decrypted_nick.unwrap();
-    if privmsg.nickname.len() > MAXIMUM_LENGTH_OF_NICKNAME {
-        privmsg.nickname = privmsg.nickname[..MAXIMUM_LENGTH_OF_NICKNAME].to_string();
-    }
     privmsg.message = decrypted_msg.unwrap();
 }
 
 /// Encrypt PrivMsg
 pub fn encrypt_privmsg(salt_box: &SalsaBox, privmsg: &mut Privmsg) {
-    privmsg.nickname = encrypt_message(salt_box, &privmsg.nickname);
-    privmsg.target = encrypt_message(salt_box, &privmsg.target);
-    privmsg.message = encrypt_message(salt_box, &privmsg.message);
+    privmsg.nickname = encrypt(salt_box, &privmsg.nickname);
+    privmsg.target = encrypt(salt_box, &privmsg.target);
+    privmsg.message = encrypt(salt_box, &privmsg.message);
 }

+ 84 - 111
bin/ircd/src/irc/client.rs

@@ -4,7 +4,7 @@ use futures::{
     io::{BufReader, ReadHalf, WriteHalf},
     AsyncBufReadExt, AsyncRead, AsyncWrite, AsyncWriteExt, FutureExt,
 };
-use fxhash::FxHashMap;
+
 use log::{debug, error, info, warn};
 
 use darkfi::{
@@ -14,12 +14,14 @@ use darkfi::{
 };
 
 use crate::{
-    buffers::{ArcPrivmsgsBuffer, SeenIds},
+    buffers::Buffers,
     crypto::{decrypt_privmsg, decrypt_target, encrypt_privmsg},
     privmsg::{MAXIMUM_LENGTH_OF_MESSAGE, MAXIMUM_LENGTH_OF_NICKNAME},
-    ChannelInfo, ContactInfo, Privmsg,
+    ChannelInfo, Privmsg,
 };
 
+use super::IrcConfig;
+
 const RPL_NOTOPIC: u32 = 331;
 const RPL_TOPIC: u32 = 332;
 const RPL_NAMEREPLY: u32 = 353;
@@ -31,25 +33,10 @@ pub struct IrcClient<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> {
     pub address: SocketAddr,
 
     // msgs buffer
-    privmsgs_buffer: ArcPrivmsgsBuffer,
-    seen_msg_ids: SeenIds,
-
-    // init bool
-    is_nick_init: bool,
-    is_user_init: bool,
-    is_registered: bool,
-    is_cap_end: bool,
-    is_pass_init: bool,
-
-    // user config
-    nickname: String,
-    password: String,
-    capabilities: FxHashMap<String, bool>,
-
-    // channels and contacts
-    auto_channels: Vec<String>,
-    pub configured_chans: FxHashMap<String, ChannelInfo>,
-    pub configured_contacts: FxHashMap<String, ContactInfo>,
+    buffers: Buffers,
+
+    // irc config
+    irc_config: IrcConfig,
 
     // p2p
     p2p: P2pPtr,
@@ -58,42 +45,16 @@ pub struct IrcClient<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> {
 }
 
 impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
-    #[allow(clippy::too_many_arguments)]
     pub fn new(
         write_stream: WriteHalf<C>,
         address: SocketAddr,
-        privmsgs_buffer: ArcPrivmsgsBuffer,
-        seen_msg_ids: SeenIds,
-        password: String,
-        auto_channels: Vec<String>,
-        configured_chans: FxHashMap<String, ChannelInfo>,
-        configured_contacts: FxHashMap<String, ContactInfo>,
+        buffers: Buffers,
+        irc_config: IrcConfig,
         p2p: P2pPtr,
         p2p_notifiers: SubscriberPtr<Privmsg>,
         subscription: Subscription<Privmsg>,
     ) -> Self {
-        let mut capabilities = FxHashMap::default();
-        capabilities.insert("no-history".to_string(), false);
-        Self {
-            write_stream,
-            address,
-            privmsgs_buffer,
-            seen_msg_ids,
-            is_nick_init: false,
-            is_user_init: false,
-            is_registered: false,
-            is_cap_end: true,
-            is_pass_init: false,
-            nickname: "anon".to_string(),
-            password,
-            auto_channels,
-            configured_chans,
-            configured_contacts,
-            capabilities,
-            p2p,
-            p2p_notifiers,
-            subscription,
-        }
+        Self { write_stream, address, buffers, irc_config, p2p, p2p_notifiers, subscription }
     }
 
     /// Start listening for messages came from p2p network or irc client
@@ -130,19 +91,21 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
         let mut msg = msg.clone();
         let mut contact = String::new();
+
         decrypt_target(
             &mut contact,
             &mut msg,
-            self.configured_chans.clone(),
-            self.configured_contacts.clone(),
+            self.irc_config.configured_chans.clone(),
+            self.irc_config.configured_contacts.clone(),
         );
+
         if msg.target.starts_with('#') {
             // Try to potentially decrypt the incoming message.
-            if !self.configured_chans.contains_key(&msg.target) {
+            if !self.irc_config.configured_chans.contains_key(&msg.target) {
                 return Ok(())
             }
 
-            let chan_info = self.configured_chans.get_mut(&msg.target).unwrap();
+            let chan_info = self.irc_config.configured_chans.get_mut(&msg.target).unwrap();
             if !chan_info.joined {
                 return Ok(())
             }
@@ -159,12 +122,12 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
             self.reply(&msg.to_string()).await?;
             return Ok(())
-        } else if self.is_cap_end && self.is_nick_init {
-            if !self.configured_contacts.contains_key(&contact) {
+        } else if self.irc_config.is_cap_end && self.irc_config.is_nick_init {
+            if !self.irc_config.configured_contacts.contains_key(&contact) {
                 return Ok(())
             }
 
-            let contact_info = self.configured_contacts.get(&contact).unwrap();
+            let contact_info = self.irc_config.configured_contacts.get(&contact).unwrap();
             if let Some(salt_box) = &contact_info.salt_box {
                 decrypt_privmsg(salt_box, &mut msg);
                 // This is for /query
@@ -201,8 +164,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             return Err(Error::MalformedPacket)
         }
 
-        if self.password.is_empty() {
-            self.is_pass_init = true
+        if self.irc_config.password.is_empty() {
+            self.irc_config.is_pass_init = true
         }
 
         let (command, value) = parse_line(&line)?;
@@ -228,23 +191,28 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
     }
 
     async fn registre(&mut self) -> Result<()> {
-        if !self.is_registered && self.is_cap_end && self.is_nick_init && self.is_user_init {
+        if !self.irc_config.is_registered &&
+            self.irc_config.is_cap_end &&
+            self.irc_config.is_nick_init &&
+            self.irc_config.is_user_init
+        {
             debug!("Initializing peer connection");
-            let register_reply = format!(":darkfi 001 {} :Let there be dark\r\n", self.nickname);
+            let register_reply =
+                format!(":darkfi 001 {} :Let there be dark\r\n", self.irc_config.nickname);
             self.reply(&register_reply).await?;
-            self.is_registered = true;
+            self.irc_config.is_registered = true;
 
-            self.on_receive_join(self.auto_channels.clone()).await?;
+            self.on_receive_join(self.irc_config.auto_channels.clone()).await?;
 
-            if *self.capabilities.get("no-history").unwrap() {
+            if *self.irc_config.capabilities.get("no-history").unwrap() {
                 return Ok(())
             }
 
             // Send dm messages in buffer
-            let privmsgs_buffer = self.privmsgs_buffer.lock().await;
+            let privmsgs_buffer = self.buffers.privmsgs.lock().await;
             for msg in privmsgs_buffer.iter() {
-                let is_dm = msg.target == self.nickname ||
-                    (msg.nickname == self.nickname && !msg.target.starts_with('#'));
+                let is_dm = msg.target == self.irc_config.nickname ||
+                    (msg.nickname == self.irc_config.nickname && !msg.target.starts_with('#'));
 
                 if is_dm {
                     self.p2p_notifiers.notify_by_id(msg.clone(), self.subscription.get_id()).await;
@@ -269,8 +237,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
     async fn on_receive_user(&mut self) -> Result<()> {
         // We can stuff any extra things like public keys in here.
         // Ignore it for now.
-        if self.is_pass_init {
-            self.is_user_init = true;
+        if self.irc_config.is_pass_init {
+            self.irc_config.is_user_init = true;
         } else {
             // Close the connection
             warn!("[CLIENT {}] Password is required", self.address);
@@ -280,8 +248,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
     }
 
     async fn on_receive_pass(&mut self, password: &str) -> Result<()> {
-        if self.password == password {
-            self.is_pass_init = true
+        if self.irc_config.password == password {
+            self.irc_config.is_pass_init = true
         } else {
             // Close the connection
             warn!("[CLIENT {}] Password is not correct!", self.address);
@@ -295,19 +263,21 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             return Ok(())
         }
 
-        self.is_nick_init = true;
-        let old_nick = std::mem::replace(&mut self.nickname, nickname.to_string());
+        self.irc_config.is_nick_init = true;
+        let old_nick = std::mem::replace(&mut self.irc_config.nickname, nickname.to_string());
 
-        let nick_reply = format!(":{}!anon@dark.fi NICK {}\r\n", old_nick, self.nickname);
+        let nick_reply =
+            format!(":{}!anon@dark.fi NICK {}\r\n", old_nick, self.irc_config.nickname);
         self.reply(&nick_reply).await
     }
 
     async fn on_receive_part(&mut self, channels: Vec<String>) -> Result<()> {
         for chan in channels.iter() {
-            let part_reply = format!(":{}!anon@dark.fi PART {}\r\n", self.nickname, chan);
+            let part_reply =
+                format!(":{}!anon@dark.fi PART {}\r\n", self.irc_config.nickname, chan);
             self.reply(&part_reply).await?;
-            if self.configured_chans.contains_key(chan) {
-                let chan_info = self.configured_chans.get_mut(chan).unwrap();
+            if self.irc_config.configured_chans.contains_key(chan) {
+                let chan_info = self.irc_config.configured_chans.get_mut(chan).unwrap();
                 chan_info.joined = false;
             }
         }
@@ -322,20 +292,22 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             }
 
             let topic = &line[substr_idx + 1..];
-            let chan_info = self.configured_chans.get_mut(channel).unwrap();
+            let chan_info = self.irc_config.configured_chans.get_mut(channel).unwrap();
             chan_info.topic = Some(topic.to_string());
 
-            let topic_reply =
-                format!(":{}!anon@dark.fi TOPIC {} :{}\r\n", self.nickname, channel, topic);
+            let topic_reply = format!(
+                ":{}!anon@dark.fi TOPIC {} :{}\r\n",
+                self.irc_config.nickname, channel, topic
+            );
             self.reply(&topic_reply).await?;
         } else {
             // Client is asking or the topic
-            let chan_info = self.configured_chans.get(channel).unwrap();
+            let chan_info = self.irc_config.configured_chans.get(channel).unwrap();
             let topic_reply = if let Some(topic) = &chan_info.topic {
-                format!("{} {} {} :{}\r\n", RPL_TOPIC, self.nickname, channel, topic)
+                format!("{} {} {} :{}\r\n", RPL_TOPIC, self.irc_config.nickname, channel, topic)
             } else {
                 const TOPIC: &str = "No topic is set";
-                format!("{} {} {} :{}\r\n", RPL_NOTOPIC, self.nickname, channel, TOPIC)
+                format!("{} {} {} :{}\r\n", RPL_NOTOPIC, self.irc_config.nickname, channel, TOPIC)
             };
             self.reply(&topic_reply).await?;
         }
@@ -348,15 +320,15 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
     }
 
     async fn on_receive_cap(&mut self, line: &str, subcommand: &str) -> Result<()> {
-        self.is_cap_end = false;
+        self.irc_config.is_cap_end = false;
 
-        let capabilities_keys: Vec<String> = self.capabilities.keys().cloned().collect();
+        let capabilities_keys: Vec<String> = self.irc_config.capabilities.keys().cloned().collect();
 
         match subcommand {
             "LS" => {
                 let cap_ls_reply = format!(
                     ":{}!anon@dark.fi CAP * LS :{}\r\n",
-                    self.nickname,
+                    self.irc_config.nickname,
                     capabilities_keys.join(" ")
                 );
                 self.reply(&cap_ls_reply).await?;
@@ -375,8 +347,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
                 let mut nak_list = vec![];
 
                 for c in cap {
-                    if self.capabilities.contains_key(c) {
-                        self.capabilities.insert(c.to_string(), true);
+                    if self.irc_config.capabilities.contains_key(c) {
+                        self.irc_config.capabilities.insert(c.to_string(), true);
                         ack_list.push(c);
                     } else {
                         nak_list.push(c);
@@ -385,13 +357,13 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
                 let cap_ack_reply = format!(
                     ":{}!anon@dark.fi CAP * ACK :{}\r\n",
-                    self.nickname,
+                    self.irc_config.nickname,
                     ack_list.join(" ")
                 );
 
                 let cap_nak_reply = format!(
                     ":{}!anon@dark.fi CAP * NAK :{}\r\n",
-                    self.nickname,
+                    self.irc_config.nickname,
                     nak_list.join(" ")
                 );
 
@@ -401,6 +373,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
             "LIST" => {
                 let enabled_capabilities: Vec<String> = self
+                    .irc_config
                     .capabilities
                     .clone()
                     .into_iter()
@@ -410,14 +383,14 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
                 let cap_list_reply = format!(
                     ":{}!anon@dark.fi CAP * LIST :{}\r\n",
-                    self.nickname,
+                    self.irc_config.nickname,
                     enabled_capabilities.join(" ")
                 );
                 self.reply(&cap_list_reply).await?;
             }
 
             "END" => {
-                self.is_cap_end = true;
+                self.irc_config.is_cap_end = true;
             }
             _ => {}
         }
@@ -429,8 +402,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             if !chan.starts_with('#') {
                 continue
             }
-            if self.configured_chans.contains_key(chan) {
-                let chan_info = self.configured_chans.get(chan).unwrap();
+            if self.irc_config.configured_chans.contains_key(chan) {
+                let chan_info = self.irc_config.configured_chans.get(chan).unwrap();
 
                 if chan_info.names.is_empty() {
                     return Ok(())
@@ -438,7 +411,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
                 let names_reply = format!(
                     ":{}!anon@dark.fi {} = {} : {}\r\n",
-                    self.nickname,
+                    self.irc_config.nickname,
                     RPL_NAMEREPLY,
                     chan,
                     chan_info.names.join(" ")
@@ -448,7 +421,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
                 let end_of_names = format!(
                     ":DarkFi {:03} {} {} :End of NAMES list\r\n",
-                    RPL_ENDOFNAMES, self.nickname, chan
+                    RPL_ENDOFNAMES, self.irc_config.nickname, chan
                 );
 
                 self.reply(&end_of_names).await?;
@@ -468,18 +441,18 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
         info!("[CLIENT {}] (Plain) PRIVMSG {} :{}", self.address, target, message,);
 
-        let privmsgs_buffer = self.privmsgs_buffer.lock().await;
+        let privmsgs_buffer = self.buffers.privmsgs.lock().await;
         let last_term = privmsgs_buffer.last_term() + 1;
         drop(privmsgs_buffer);
 
-        let mut privmsg = Privmsg::new(&self.nickname, target, &message, last_term);
+        let mut privmsg = Privmsg::new(&self.irc_config.nickname, target, &message, last_term);
 
         if target.starts_with('#') {
-            if !self.configured_chans.contains_key(target) {
+            if !self.irc_config.configured_chans.contains_key(target) {
                 return Ok(())
             }
 
-            let channel_info = self.configured_chans.get(target).unwrap();
+            let channel_info = self.irc_config.configured_chans.get(target).unwrap();
 
             if !channel_info.joined {
                 return Ok(())
@@ -490,11 +463,11 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
                 info!("[CLIENT {}] (Encrypted) PRIVMSG: {:?}", self.address, privmsg);
             }
         } else {
-            if !self.configured_contacts.contains_key(target) {
+            if !self.irc_config.configured_contacts.contains_key(target) {
                 return Ok(())
             }
 
-            let contact_info = self.configured_contacts.get(target).unwrap();
+            let contact_info = self.irc_config.configured_contacts.get(target).unwrap();
             if let Some(salt_box) = &contact_info.salt_box {
                 encrypt_privmsg(salt_box, &mut privmsg);
                 info!("[CLIENT {}] (Encrypted) PRIVMSG: {:?}", self.address, privmsg);
@@ -502,8 +475,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
         }
 
         {
-            (*self.seen_msg_ids.lock().await).push(privmsg.id);
-            (*self.privmsgs_buffer.lock().await).push(&privmsg)
+            (*self.buffers.seen_ids.lock().await).push(privmsg.id);
+            (*self.buffers.unread_msgs.lock().await).insert(&privmsg);
         }
 
         self.p2p_notifiers
@@ -521,13 +494,13 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             if !chan.starts_with('#') {
                 continue
             }
-            if !self.configured_chans.contains_key(chan) {
+            if !self.irc_config.configured_chans.contains_key(chan) {
                 let mut chan_info = ChannelInfo::new()?;
                 chan_info.topic = Some("n/a".to_string());
-                self.configured_chans.insert(chan.to_string(), chan_info);
+                self.irc_config.configured_chans.insert(chan.to_string(), chan_info);
             }
 
-            let chan_info = self.configured_chans.get_mut(chan).unwrap();
+            let chan_info = self.irc_config.configured_chans.get_mut(chan).unwrap();
             if chan_info.joined {
                 return Ok(())
             }
@@ -538,15 +511,15 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             chan_info.topic = Some(topic.to_string());
 
             {
-                let j = format!(":{}!anon@dark.fi JOIN {}\r\n", self.nickname, chan);
+                let j = format!(":{}!anon@dark.fi JOIN {}\r\n", self.irc_config.nickname, chan);
                 let t = format!(":DarkFi TOPIC {} :{}\r\n", chan, topic);
                 self.reply(&j).await?;
                 self.reply(&t).await?;
             }
 
             // Send messages in buffer
-            if !self.capabilities.get("no-history").unwrap() {
-                for msg in self.privmsgs_buffer.lock().await.iter() {
+            if !self.irc_config.capabilities.get("no-history").unwrap() {
+                for msg in self.buffers.privmsgs.lock().await.iter() {
                     if msg.target == *chan {
                         self.p2p_notifiers
                             .notify_by_id(msg.clone(), self.subscription.get_id())

+ 69 - 33
bin/ircd/src/irc/mod.rs

@@ -7,26 +7,80 @@ use futures_rustls::{rustls, TlsAcceptor};
 use fxhash::FxHashMap;
 use log::{error, info};
 
-use darkfi::{net::P2pPtr, system::SubscriberPtr, util::expand_path, Error, Result};
+use darkfi::{
+    net::P2pPtr,
+    system::SubscriberPtr,
+    util::{expand_path, path::get_config_path},
+    Error, Result,
+};
 
 use crate::{
-    buffers::{ArcPrivmsgsBuffer, SeenIds},
-    settings::Args,
-    ChannelInfo, ContactInfo, Privmsg,
+    buffers::Buffers,
+    settings::{
+        parse_configured_channels, parse_configured_contacts, Args, ChannelInfo, ContactInfo,
+        CONFIG_FILE,
+    },
+    Privmsg,
 };
 
 mod client;
 
 pub use client::IrcClient;
 
+#[derive(Clone)]
+pub struct IrcConfig {
+    // init bool
+    pub is_nick_init: bool,
+    pub is_user_init: bool,
+    pub is_registered: bool,
+    pub is_cap_end: bool,
+    pub is_pass_init: bool,
+
+    // user config
+    pub nickname: String,
+    pub password: String,
+    pub capabilities: FxHashMap<String, bool>,
+
+    // channels and contacts
+    pub auto_channels: Vec<String>,
+    pub configured_chans: FxHashMap<String, ChannelInfo>,
+    pub configured_contacts: FxHashMap<String, ContactInfo>,
+}
+
+impl IrcConfig {
+    pub fn new(settings: &Args) -> Result<Self> {
+        let password = settings.password.as_ref().unwrap_or(&String::new()).clone();
+
+        let auto_channels = settings.autojoin.clone();
+
+        // Pick up channel settings from the TOML configuration
+        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)?;
+
+        let mut capabilities = FxHashMap::default();
+        capabilities.insert("no-history".to_string(), false);
+        Ok(Self {
+            is_nick_init: false,
+            is_user_init: false,
+            is_registered: false,
+            is_cap_end: true,
+            is_pass_init: false,
+            nickname: "anon".to_string(),
+            password,
+            auto_channels,
+            configured_chans,
+            configured_contacts,
+            capabilities,
+        })
+    }
+}
+
 pub struct IrcServer {
     settings: Args,
-    privmsgs_buffer: ArcPrivmsgsBuffer,
-    seen_msg_ids: SeenIds,
-    auto_channels: Vec<String>,
-    password: String,
-    configured_chans: FxHashMap<String, ChannelInfo>,
-    configured_contacts: FxHashMap<String, ContactInfo>,
+    irc_config: IrcConfig,
+    buffers: Buffers,
     p2p: P2pPtr,
     p2p_notifiers: SubscriberPtr<Privmsg>,
 }
@@ -34,26 +88,12 @@ pub struct IrcServer {
 impl IrcServer {
     pub async fn new(
         settings: Args,
-        privmsgs_buffer: ArcPrivmsgsBuffer,
-        seen_msg_ids: SeenIds,
-        auto_channels: Vec<String>,
-        password: String,
-        configured_chans: FxHashMap<String, ChannelInfo>,
-        configured_contacts: FxHashMap<String, ContactInfo>,
+        buffers: Buffers,
         p2p: P2pPtr,
         p2p_notifiers: SubscriberPtr<Privmsg>,
     ) -> Result<Self> {
-        Ok(Self {
-            settings,
-            privmsgs_buffer,
-            seen_msg_ids,
-            auto_channels,
-            password,
-            configured_chans,
-            configured_contacts,
-            p2p,
-            p2p_notifiers,
-        })
+        let irc_config = IrcConfig::new(&settings)?;
+        Ok(Self { settings, irc_config, buffers, p2p, p2p_notifiers })
     }
 
     /// Start listening to new irc clients connecting to the irc server address
@@ -111,12 +151,8 @@ impl IrcServer {
         let mut client = IrcClient::new(
             writer,
             peer_addr,
-            self.privmsgs_buffer.clone(),
-            self.seen_msg_ids.clone(),
-            self.password.clone(),
-            self.auto_channels.clone(),
-            self.configured_chans.clone(),
-            self.configured_contacts.clone(),
+            self.buffers.clone(),
+            self.irc_config.clone(),
             self.p2p.clone(),
             self.p2p_notifiers.clone(),
             p2p_subscription,

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

@@ -4,7 +4,6 @@ use std::fmt;
 use async_channel::Receiver;
 use async_executor::Executor;
 
-use fxhash::FxHashMap;
 use log::{info, warn};
 use rand::rngs::OsRng;
 use smol::future;
@@ -32,15 +31,12 @@ pub mod rpc;
 pub mod settings;
 
 use crate::{
-    buffers::{ArcPrivmsgsBuffer, PrivmsgsBuffer, RingBuffer, SeenIds, SIZE_OF_MSG_IDSS_BUFFER},
+    buffers::{create_buffers, Buffers, RingBuffer, SIZE_OF_MSG_IDSS_BUFFER},
     irc::IrcServer,
     privmsg::Privmsg,
     protocol_privmsg::ProtocolPrivmsg,
     rpc::JsonRpcInterface,
-    settings::{
-        parse_configured_channels, parse_configured_contacts, Args, ChannelInfo, ContactInfo,
-        CONFIG_FILE, CONFIG_FILE_CONTENTS,
-    },
+    settings::{Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
 };
 
 #[derive(serde::Serialize)]
@@ -55,50 +51,23 @@ impl fmt::Display for KeyPair {
     }
 }
 
-pub type UnreadMsgs = Arc<Mutex<FxHashMap<String, Privmsg>>>;
-
 struct Ircd {
-    // msgs
-    privmsgs_buffer: ArcPrivmsgsBuffer,
-    seen_msg_ids: SeenIds,
-    // channels
-    autojoin_chans: Vec<String>,
-    configured_chans: FxHashMap<String, ChannelInfo>,
-    configured_contacts: FxHashMap<String, ContactInfo>,
-    // p2p
-    p2p: net::P2pPtr,
     p2p_notifiers: SubscriberPtr<Privmsg>,
-    password: String,
 }
 
 impl Ircd {
-    fn new(
-        privmsgs_buffer: ArcPrivmsgsBuffer,
-        seen_msg_ids: SeenIds,
-        autojoin_chans: Vec<String>,
-        password: String,
-        configured_chans: FxHashMap<String, ChannelInfo>,
-        configured_contacts: FxHashMap<String, ContactInfo>,
-        p2p: net::P2pPtr,
-    ) -> Self {
+    fn new() -> Self {
         let p2p_notifiers = Subscriber::new();
-        Self {
-            privmsgs_buffer,
-            seen_msg_ids,
-            autojoin_chans,
-            password,
-            configured_chans,
-            configured_contacts,
-            p2p,
-            p2p_notifiers,
-        }
+        Self { p2p_notifiers }
     }
 
     async fn start(
         &self,
         settings: &Args,
-        executor: Arc<Executor<'_>>,
+        buffers: Buffers,
+        p2p: net::P2pPtr,
         p2p_receiver: Receiver<Privmsg>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<()> {
         let p2p_notifiers = self.p2p_notifiers.clone();
         executor
@@ -111,13 +80,8 @@ impl Ircd {
 
         let irc_server = IrcServer::new(
             settings.clone(),
-            self.privmsgs_buffer.clone(),
-            self.seen_msg_ids.clone(),
-            self.autojoin_chans.clone(),
-            self.password.clone(),
-            self.configured_chans.clone(),
-            self.configured_contacts.clone(),
-            self.p2p.clone(),
+            buffers.clone(),
+            p2p.clone(),
             self.p2p_notifiers.clone(),
         )
         .await?;
@@ -129,10 +93,8 @@ impl Ircd {
 
 async_daemonize!(realmain);
 async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
-    let seen_msg_ids = Arc::new(Mutex::new(RingBuffer::new(SIZE_OF_MSG_IDSS_BUFFER)));
     let seen_inv_ids = Arc::new(Mutex::new(RingBuffer::new(SIZE_OF_MSG_IDSS_BUFFER)));
-    let privmsgs_buffer = PrivmsgsBuffer::new();
-    let unread_msgs = Arc::new(Mutex::new(FxHashMap::default()));
+    let buffers = create_buffers();
 
     if settings.gen_secret {
         let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
@@ -159,14 +121,6 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
         return Ok(())
     }
 
-    let password = settings.password.clone().unwrap_or_default();
-
-    // Pick up channel settings from the TOML configuration
-    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)?;
-
     //
     // P2p setup
     //
@@ -179,28 +133,16 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
 
     let registry = p2p.protocol_registry();
 
-    let seen_msg_ids_cloned = seen_msg_ids.clone();
+    let buffers_cloned = buffers.clone();
     let seen_inv_ids_cloned = seen_inv_ids.clone();
-    let privmsgs_buffer_cloned = privmsgs_buffer.clone();
-    let unread_msgs_cloned = unread_msgs.clone();
     registry
         .register(net::SESSION_ALL, move |channel, p2p| {
             let sender = p2p_send_channel.clone();
-            let seen_msg_ids_cloned = seen_msg_ids_cloned.clone();
             let seen_inv_ids_cloned = seen_inv_ids_cloned.clone();
-            let privmsgs_buffer_cloned = privmsgs_buffer_cloned.clone();
-            let unread_msgs_cloned = unread_msgs_cloned.clone();
+            let buffers_cloned = buffers_cloned.clone();
             async move {
-                ProtocolPrivmsg::init(
-                    channel,
-                    sender,
-                    p2p,
-                    seen_msg_ids_cloned,
-                    seen_inv_ids_cloned,
-                    privmsgs_buffer_cloned,
-                    unread_msgs_cloned,
-                )
-                .await
+                ProtocolPrivmsg::init(channel, sender, p2p, seen_inv_ids_cloned, buffers_cloned)
+                    .await
             }
         })
         .await;
@@ -224,17 +166,9 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     // IRC instance
     //
 
-    let ircd = Ircd::new(
-        privmsgs_buffer.clone(),
-        seen_msg_ids.clone(),
-        settings.autojoin.clone(),
-        password.clone(),
-        configured_chans.clone(),
-        configured_contacts.clone(),
-        p2p.clone(),
-    );
-
-    ircd.start(&settings, executor.clone(), p2p_recv_channel).await?;
+    let ircd = Ircd::new();
+
+    ircd.start(&settings, buffers, p2p, p2p_recv_channel, executor.clone()).await?;
 
     // Run once receive exit signal
     let (signal, shutdown) = async_channel::bounded::<()>(1);

+ 16 - 28
bin/ircd/src/protocol_privmsg.rs

@@ -5,7 +5,6 @@ use async_trait::async_trait;
 use chrono::Utc;
 use log::debug;
 use rand::{rngs::OsRng, RngCore};
-use ripemd::{Digest, Ripemd160};
 
 use darkfi::{
     net,
@@ -17,8 +16,8 @@ use darkfi::{
 };
 
 use crate::{
-    buffers::{ArcPrivmsgsBuffer, SeenIds},
-    Privmsg, UnreadMsgs,
+    buffers::{Buffers, InvSeenIds},
+    Privmsg,
 };
 
 const MAX_CONFIRM: u8 = 4;
@@ -59,11 +58,9 @@ pub struct ProtocolPrivmsg {
     inv_sub: net::MessageSubscription<Inv>,
     getdata_sub: net::MessageSubscription<GetData>,
     p2p: net::P2pPtr,
-    msg_ids: SeenIds,
-    inv_ids: SeenIds,
-    msgs: ArcPrivmsgsBuffer,
-    unread_msgs: UnreadMsgs,
     channel: net::ChannelPtr,
+    inv_ids: InvSeenIds,
+    buffers: Buffers,
 }
 
 impl ProtocolPrivmsg {
@@ -71,10 +68,8 @@ impl ProtocolPrivmsg {
         channel: net::ChannelPtr,
         notify: async_channel::Sender<Privmsg>,
         p2p: net::P2pPtr,
-        msg_ids: SeenIds,
-        inv_ids: SeenIds,
-        msgs: ArcPrivmsgsBuffer,
-        unread_msgs: UnreadMsgs,
+        inv_ids: InvSeenIds,
+        buffers: Buffers,
     ) -> net::ProtocolBasePtr {
         let message_subsytem = channel.get_message_subsystem();
         message_subsytem.add_dispatch::<Privmsg>().await;
@@ -96,11 +91,9 @@ impl ProtocolPrivmsg {
             getdata_sub,
             jobsman: net::ProtocolJobsManager::new("ProtocolPrivmsg", channel.clone()),
             p2p,
-            msg_ids,
-            inv_ids,
-            msgs,
-            unread_msgs,
             channel,
+            inv_ids,
+            buffers,
         })
     }
 
@@ -120,7 +113,7 @@ impl ProtocolPrivmsg {
 
             let mut inv_requested = vec![];
             for inv_object in inv.invs.iter() {
-                let mut msgs = self.unread_msgs.lock().await;
+                let msgs = &mut self.buffers.unread_msgs.lock().await.0;
                 if let Some(msg) = msgs.get_mut(&inv_object.0) {
                     msg.read_confirms += 1;
                 } else {
@@ -145,7 +138,7 @@ impl ProtocolPrivmsg {
             let msg = self.msg_sub.receive().await?;
             let mut msg = (*msg).to_owned();
 
-            let mut msg_ids = self.msg_ids.lock().await;
+            let mut msg_ids = self.buffers.seen_ids.lock().await;
             if msg_ids.contains(&msg.id) {
                 continue
             }
@@ -171,7 +164,7 @@ impl ProtocolPrivmsg {
             let getdata = self.getdata_sub.receive().await?;
             let getdata = (*getdata).to_owned();
 
-            let msgs = self.unread_msgs.lock().await;
+            let msgs = &self.buffers.unread_msgs.lock().await.0;
             for inv in getdata.invs {
                 if let Some(msg) = msgs.get(&inv.0) {
                     self.channel.send(msg.clone()).await?;
@@ -181,16 +174,11 @@ impl ProtocolPrivmsg {
     }
 
     async fn add_to_unread_msgs(&self, msg: &Privmsg) -> String {
-        let mut msgs = self.unread_msgs.lock().await;
-        let mut hasher = Ripemd160::new();
-        hasher.update(msg.to_string());
-        let key = hex::encode(hasher.finalize());
-        msgs.insert(key.clone(), msg.clone());
-        key
+        self.buffers.unread_msgs.lock().await.insert(msg)
     }
 
     async fn update_unread_msgs(&self) -> Result<()> {
-        let mut msgs = self.unread_msgs.lock().await;
+        let msgs = &mut self.buffers.unread_msgs.lock().await.0;
         for (hash, msg) in msgs.clone() {
             if msg.timestamp + UNREAD_MSG_EXPIRE_TIME < Utc::now().timestamp() {
                 msgs.remove(&hash);
@@ -205,7 +193,7 @@ impl ProtocolPrivmsg {
     }
 
     async fn add_to_msgs(&self, msg: &Privmsg) -> Result<()> {
-        self.msgs.lock().await.push(msg);
+        self.buffers.privmsgs.lock().await.push(msg);
         self.notify.send(msg.clone()).await?;
         Ok(())
     }
@@ -215,7 +203,7 @@ impl ProtocolPrivmsg {
 
         self.update_unread_msgs().await?;
 
-        for msg in self.unread_msgs.lock().await.values() {
+        for msg in self.buffers.unread_msgs.lock().await.0.values() {
             self.channel.send(msg.clone()).await?;
         }
         Ok(())
@@ -229,7 +217,7 @@ impl net::ProtocolBase for ProtocolPrivmsg {
     /// waits for pong reply. Waits for ping and replies with a pong.
     async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
         // once a channel get started
-        let msgs_buffer = self.msgs.lock().await;
+        let msgs_buffer = self.buffers.privmsgs.lock().await;
         for m in msgs_buffer.iter() {
             self.channel.send(m.clone()).await?;
         }