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

bin/ircd: use ring buffer for messages history

ghassmo 4 лет назад
Родитель
Сommit
40e20ba078
5 измененных файлов с 55 добавлено и 24 удалено
  1. 16 0
      Cargo.lock
  2. 1 0
      bin/ircd/Cargo.toml
  3. 25 20
      bin/ircd/src/main.rs
  4. 4 0
      bin/ircd/src/privmsg.rs
  5. 9 4
      bin/ircd/src/server.rs

+ 16 - 0
Cargo.lock

@@ -68,6 +68,12 @@ version = "1.0.57"
 source = "registry+https://github.com/rust-lang/crates.io-index"
 checksum = "08f9b8508dccb7687a1d6c4ce66b2b0ecef467c94667de27d8d7fe1f8d2a9cdc"
 
+[[package]]
+name = "array-init"
+version = "2.0.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "6945cc5422176fc5e602e590c2878d2c2acd9a4fe20a4baa7c28022521698ec6"
+
 [[package]]
 name = "arrayref"
 version = "0.3.6"
@@ -2280,6 +2286,7 @@ dependencies = [
  "fxhash",
  "log",
  "rand",
+ "ringbuffer",
  "serde",
  "serde_json",
  "simplelog",
@@ -3313,6 +3320,15 @@ dependencies = [
  "winapi",
 ]
 
+[[package]]
+name = "ringbuffer"
+version = "0.8.4"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "b30a00730a27595dcf899dce512aa031dd650f86aafcb132fd8dd9f409b369d0"
+dependencies = [
+ "array-init",
+]
+
 [[package]]
 name = "rkyv"
 version = "0.7.38"

+ 1 - 0
bin/ircd/Cargo.toml

@@ -31,6 +31,7 @@ simplelog = "0.12.0"
 fxhash = "0.2.1"
 ctrlc-async = {version= "3.2.2", default-features = false, features = ["async-std", "termination"]}
 url = "2.2.2"
+ringbuffer = "0.8.4" 
 
 # Encoding and parsing
 serde_json = "1.0.81"

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

@@ -11,6 +11,7 @@ use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
 use fxhash::FxHashMap;
 use log::{debug, error, info, warn};
 use rand::rngs::OsRng;
+use ringbuffer::RingBufferWrite;
 use smol::future;
 use structopt_toml::StructOptToml;
 
@@ -33,13 +34,15 @@ pub mod settings;
 
 use crate::{
     crypto::try_decrypt_message,
-    privmsg::{Privmsg, SeenMsgIds},
+    privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
     protocol_privmsg::ProtocolPrivmsg,
     rpc::JsonRpcInterface,
     server::IrcServerConnection,
     settings::{parse_configured_channels, Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
 };
 
+const SIZE_OF_MSGS_BUFFER: usize = 4096;
+
 fn build_irc_msg(msg: &Privmsg) -> String {
     debug!("ABOUT TO SEND: {:?}", msg);
     let irc_msg =
@@ -65,27 +68,13 @@ fn clean_input(mut line: String, peer_addr: &SocketAddr) -> Result<String> {
     Ok(line)
 }
 
-async fn broadcast_msg(
-    irc_msg: String,
-    peer_addr: SocketAddr,
-    conn: &mut IrcServerConnection,
-) -> Result<()> {
-    info!("Send msg to IRC client '{}' from {}", irc_msg, peer_addr);
-
-    if let Err(e) = conn.update(irc_msg).await {
-        warn!("Connection error: {} for {}", e, peer_addr);
-        return Err(Error::ChannelStopped)
-    }
-
-    Ok(())
-}
-
 async fn process(
     // server
     stream: TcpStream,
     peer_addr: SocketAddr,
-    // msg ids
+    // msgs
     seen_msg_ids: SeenMsgIds,
+    privmsgs_buffer: PrivmsgsBuffer,
     // channels
     autojoin_chans: Vec<String>,
     configured_chans: FxHashMap<String, ChannelInfo>,
@@ -99,6 +88,7 @@ async fn process(
     let mut conn = IrcServerConnection::new(
         writer,
         seen_msg_ids.clone(),
+        privmsgs_buffer.clone(),
         autojoin_chans,
         configured_chans,
         p2p.clone(),
@@ -109,7 +99,7 @@ async fn process(
         futures::select! {
             privmsg = p2p_receiver.recv().fuse() => {
                 let mut msg = privmsg?;
-                info!("Received msg from Raft: {:?}", msg);
+                info!("Received msg from P2p network: {:?}", msg);
 
                 // Try to potentially decrypt the incoming message.
                 if conn.configured_chans.contains_key(&msg.channel) {
@@ -125,6 +115,11 @@ async fn process(
                     }
                 }
 
+                // add the msg to buffer
+                {
+                    (*privmsgs_buffer.lock().await).push(msg.clone());
+                }
+
                 let irc_msg = build_irc_msg(&msg);
                 conn.reply(&irc_msg).await?;
             }
@@ -138,7 +133,13 @@ async fn process(
                     Ok(m) => m,
                     Err(e) => return Err(e)
                 };
-                broadcast_msg(irc_msg, peer_addr,&mut conn).await?;
+
+                info!("Send msg to IRC client '{}' from {}", irc_msg, peer_addr);
+
+                if let Err(e) = conn.update(irc_msg).await {
+                    warn!("Connection error: {} for {}", e, peer_addr);
+                    return Err(Error::ChannelStopped)
+                }
             }
         };
     }
@@ -146,6 +147,10 @@ async fn process(
 
 async_daemonize!(realmain);
 async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
+    let seen_msg_ids = Arc::new(Mutex::new(vec![]));
+    let privmsgs_buffer: PrivmsgsBuffer =
+        Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSGS_BUFFER)));
+
     if settings.gen_secret {
         let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
         let encoded = bs58::encode(secret_key.as_bytes());
@@ -168,7 +173,6 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
 
     let registry = p2p.protocol_registry();
 
-    let seen_msg_ids = Arc::new(Mutex::new(vec![]));
     let seen_msg_ids_cloned = seen_msg_ids.clone();
     registry
         .register(net::SESSION_ALL, move |channel, p2p| {
@@ -218,6 +222,7 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
                         stream,
                         peer_addr,
                         seen_msg_ids.clone(),
+                        privmsgs_buffer.clone(),
                         settings.autojoin.clone(),
                         configured_chans.clone(),
                         p2p.clone(),

+ 4 - 0
bin/ircd/src/privmsg.rs

@@ -1,11 +1,15 @@
 use async_std::sync::{Arc, Mutex};
 
+use ringbuffer::AllocRingBuffer;
+
 use darkfi::util::serial::{SerialDecodable, SerialEncodable};
 
 pub type PrivmsgId = u64;
 
 pub type SeenMsgIds = Arc<Mutex<Vec<u64>>>;
 
+pub type PrivmsgsBuffer = Arc<Mutex<AllocRingBuffer<Privmsg>>>;
+
 #[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
 pub struct Privmsg {
     pub id: PrivmsgId,

+ 9 - 4
bin/ircd/src/server.rs

@@ -5,12 +5,13 @@ use futures::{io::WriteHalf, AsyncWriteExt};
 use fxhash::FxHashMap;
 use log::{debug, info, warn};
 use rand::{rngs::OsRng, RngCore};
+use ringbuffer::RingBufferWrite;
 
 use darkfi::{net::P2pPtr, Error, Result};
 
 use crate::{
     crypto::encrypt_message,
-    privmsg::{Privmsg, SeenMsgIds},
+    privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
     ChannelInfo,
 };
 
@@ -22,6 +23,7 @@ pub struct IrcServerConnection {
     write_stream: WriteHalf<TcpStream>,
     // msg ids
     seen_msg_ids: SeenMsgIds,
+    privmsgs_buffer: PrivmsgsBuffer,
     // user & channels
     is_nick_init: bool,
     is_user_init: bool,
@@ -37,6 +39,7 @@ impl IrcServerConnection {
     pub fn new(
         write_stream: WriteHalf<TcpStream>,
         seen_msg_ids: SeenMsgIds,
+        privmsgs_buffer: PrivmsgsBuffer,
         auto_channels: Vec<String>,
         configured_chans: FxHashMap<String, ChannelInfo>,
         p2p: P2pPtr,
@@ -44,6 +47,7 @@ impl IrcServerConnection {
         Self {
             write_stream,
             seen_msg_ids,
+            privmsgs_buffer,
             is_nick_init: false,
             is_user_init: false,
             is_registered: false,
@@ -170,9 +174,10 @@ impl IrcServerConnection {
                             message,
                         };
 
-                        let mut smi = self.seen_msg_ids.lock().await;
-                        smi.push(random_id);
-                        drop(smi);
+                        {
+                            (*self.seen_msg_ids.lock().await).push(random_id);
+                            (*self.privmsgs_buffer.lock().await).push(protocol_msg.clone())
+                        }
 
                         debug!(target: "ircd", "PRIVMSG to be sent: {:?}", protocol_msg);
                         self.p2p.broadcast(protocol_msg).await?;