Parcourir la source

bin/ircd: add timestamp to privmsg

ghassmo il y a 4 ans
Parent
commit
8d7238968b
4 fichiers modifiés avec 49 ajouts et 19 suppressions
  1. 5 6
      bin/ircd/src/main.rs
  2. 33 3
      bin/ircd/src/privmsg.rs
  3. 4 4
      bin/ircd/src/protocol_privmsg.rs
  4. 7 6
      bin/ircd/src/server.rs

+ 5 - 6
bin/ircd/src/main.rs

@@ -34,7 +34,7 @@ pub mod server;
 pub mod settings;
 
 use crate::{
-    privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
+    privmsg::{ArcPrivmsgsBuffer, Privmsg, PrivmsgsBuffer, SeenMsgIds},
     protocol_privmsg::ProtocolPrivmsg,
     rpc::JsonRpcInterface,
     server::IrcServerConnection,
@@ -45,14 +45,14 @@ use crate::{
 };
 
 const SIZE_OF_MSG_IDSS_BUFFER: usize = 65536;
-const SIZE_OF_MSGS_BUFFER: usize = 4096;
+pub const SIZE_OF_MSGS_BUFFER: usize = 4096;
 pub const MAXIMUM_LENGTH_OF_MESSAGE: usize = 1024;
 pub const MAXIMUM_LENGTH_OF_NICKNAME: usize = 32;
 
 struct Ircd {
     // msgs
     seen_msg_ids: SeenMsgIds,
-    privmsgs_buffer: PrivmsgsBuffer,
+    privmsgs_buffer: ArcPrivmsgsBuffer,
     // channels
     autojoin_chans: Vec<String>,
     configured_chans: FxHashMap<String, ChannelInfo>,
@@ -66,7 +66,7 @@ struct Ircd {
 impl Ircd {
     fn new(
         seen_msg_ids: SeenMsgIds,
-        privmsgs_buffer: PrivmsgsBuffer,
+        privmsgs_buffer: ArcPrivmsgsBuffer,
         autojoin_chans: Vec<String>,
         password: String,
         configured_chans: FxHashMap<String, ChannelInfo>,
@@ -168,8 +168,7 @@ async_daemonize!(realmain);
 async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
     let seen_msg_ids =
         Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSG_IDSS_BUFFER)));
-    let privmsgs_buffer: PrivmsgsBuffer =
-        Arc::new(Mutex::new(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSGS_BUFFER)));
+    let privmsgs_buffer = PrivmsgsBuffer::new();
 
     if settings.gen_secret {
         let secret_key = crypto_box::SecretKey::generate(&mut OsRng);

+ 33 - 3
bin/ircd/src/privmsg.rs

@@ -1,14 +1,43 @@
 use async_std::sync::{Arc, Mutex};
 
-use ringbuffer::AllocRingBuffer;
+use ringbuffer::{AllocRingBuffer, RingBufferExt, RingBufferWrite};
 
-use darkfi::util::serial::{SerialDecodable, SerialEncodable};
+use darkfi::util::{
+    serial::{SerialDecodable, SerialEncodable},
+    Timestamp,
+};
+
+use crate::SIZE_OF_MSGS_BUFFER;
 
 pub type PrivmsgId = u64;
 
 pub type SeenMsgIds = Arc<Mutex<AllocRingBuffer<u64>>>;
 
-pub type PrivmsgsBuffer = Arc<Mutex<AllocRingBuffer<Privmsg>>>;
+pub type ArcPrivmsgsBuffer = Arc<Mutex<PrivmsgsBuffer>>;
+
+pub struct PrivmsgsBuffer(AllocRingBuffer<Privmsg>);
+
+impl PrivmsgsBuffer {
+    pub fn new() -> ArcPrivmsgsBuffer {
+        Arc::new(Mutex::new(Self(ringbuffer::AllocRingBuffer::with_capacity(SIZE_OF_MSGS_BUFFER))))
+    }
+
+    pub fn push(&mut self, privmsg: &Privmsg) {
+        if privmsg.timestamp > Timestamp::current_time() {
+            return
+        }
+
+        if let Some(last_msg) = self.0.get(-1) {
+            if privmsg.timestamp > last_msg.timestamp {
+                self.0.push(privmsg.clone());
+            }
+        }
+    }
+
+    pub fn to_vec(&self) -> Vec<Privmsg> {
+        self.0.to_vec().clone()
+    }
+}
 
 #[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
 pub struct Privmsg {
@@ -16,6 +45,7 @@ pub struct Privmsg {
     pub nickname: String,
     pub target: String,
     pub message: String,
+    pub timestamp: Timestamp,
 }
 
 impl Privmsg {

+ 4 - 4
bin/ircd/src/protocol_privmsg.rs

@@ -7,7 +7,7 @@ use ringbuffer::{RingBufferExt, RingBufferWrite};
 
 use darkfi::{net, Result};
 
-use crate::privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds};
+use crate::privmsg::{ArcPrivmsgsBuffer, Privmsg, SeenMsgIds};
 
 pub struct ProtocolPrivmsg {
     jobsman: net::ProtocolJobsManagerPtr,
@@ -15,7 +15,7 @@ pub struct ProtocolPrivmsg {
     msg_sub: net::MessageSubscription<Privmsg>,
     p2p: net::P2pPtr,
     msg_ids: SeenMsgIds,
-    msgs: PrivmsgsBuffer,
+    msgs: ArcPrivmsgsBuffer,
     channel: net::ChannelPtr,
 }
 
@@ -25,7 +25,7 @@ impl ProtocolPrivmsg {
         notify_queue_sender: async_channel::Sender<Privmsg>,
         p2p: net::P2pPtr,
         msg_ids: SeenMsgIds,
-        msgs: PrivmsgsBuffer,
+        msgs: ArcPrivmsgsBuffer,
     ) -> net::ProtocolBasePtr {
         let message_subsytem = channel.get_message_subsystem();
         message_subsytem.add_dispatch::<Privmsg>().await;
@@ -70,7 +70,7 @@ impl ProtocolPrivmsg {
             }
 
             // add the msg to the buffer
-            self.msgs.lock().await.push(msg.clone());
+            self.msgs.lock().await.push(&msg);
 
             self.notify_queue_sender.send(msg.clone()).await?;
 

+ 7 - 6
bin/ircd/src/server.rs

@@ -4,13 +4,13 @@ use futures::{io::WriteHalf, AsyncRead, AsyncWrite, AsyncWriteExt};
 use fxhash::FxHashMap;
 use log::{debug, info, warn};
 use rand::{rngs::OsRng, RngCore};
-use ringbuffer::{RingBufferExt, RingBufferWrite};
+use ringbuffer::RingBufferWrite;
 
-use darkfi::{net::P2pPtr, system::SubscriberPtr, Error, Result};
+use darkfi::{net::P2pPtr, system::SubscriberPtr, util::Timestamp, Error, Result};
 
 use crate::{
     crypto::{decrypt_privmsg, decrypt_target, encrypt_privmsg},
-    privmsg::{Privmsg, PrivmsgsBuffer, SeenMsgIds},
+    privmsg::{ArcPrivmsgsBuffer, Privmsg, SeenMsgIds},
     ChannelInfo, MAXIMUM_LENGTH_OF_MESSAGE, MAXIMUM_LENGTH_OF_NICKNAME,
 };
 
@@ -25,7 +25,7 @@ pub struct IrcServerConnection<C: AsyncRead + AsyncWrite + Send + Unpin + 'stati
     peer_address: SocketAddr,
     // msg ids
     seen_msg_ids: SeenMsgIds,
-    privmsgs_buffer: PrivmsgsBuffer,
+    privmsgs_buffer: ArcPrivmsgsBuffer,
     // user & channels
     is_nick_init: bool,
     is_user_init: bool,
@@ -50,7 +50,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcServerConnection<C>
         write_stream: WriteHalf<C>,
         peer_address: SocketAddr,
         seen_msg_ids: SeenMsgIds,
-        privmsgs_buffer: PrivmsgsBuffer,
+        privmsgs_buffer: ArcPrivmsgsBuffer,
         auto_channels: Vec<String>,
         password: String,
         configured_chans: FxHashMap<String, ChannelInfo>,
@@ -217,6 +217,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcServerConnection<C>
                     nickname: self.nickname.clone(),
                     target: target.to_string().clone(),
                     message,
+                    timestamp: Timestamp::current_time(),
                 };
 
                 if target.starts_with('#') {
@@ -393,7 +394,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcServerConnection<C>
     async fn on_receive_privmsg(&mut self, privmsg: Privmsg) -> Result<()> {
         {
             (*self.seen_msg_ids.lock().await).push(privmsg.id);
-            (*self.privmsgs_buffer.lock().await).push(privmsg.clone())
+            (*self.privmsgs_buffer.lock().await).push(&privmsg)
         }
 
         self.senders.notify_with_exclude(privmsg.clone(), &[self.subscriber_id]).await;