Quellcode durchsuchen

src/system: remove notify_by_id function from subscriber

ghassmo vor 3 Jahren
Ursprung
Commit
cd200f8be4
3 geänderte Dateien mit 7 neuen und 35 gelöschten Zeilen
  1. 1 0
      bin/ircd/src/buffers.rs
  2. 6 24
      bin/ircd/src/irc/client.rs
  3. 0 11
      src/system/subscriber.rs

+ 1 - 0
bin/ircd/src/buffers.rs

@@ -100,6 +100,7 @@ impl<T: Eq + PartialEq + Clone> RingBuffer<T> {
     }
 }
 
+#[derive(Clone)]
 pub struct PrivmsgsBuffer {
     buffer: RingBuffer<Privmsg>,
     orphans: RingBuffer<Orphan>,

+ 6 - 24
bin/ircd/src/irc/client.rs

@@ -64,8 +64,8 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
 
             futures::select! {
                 msg = self.subscription.receive().fuse() => {
-                    if let Err(e) = self.process_msg_from_p2p(&msg).await {
-                        error!("[CLIENT {}] Process msg from p2p: {}",  self.address, e);
+                    if let Err(e) = self.process_msg(&msg).await {
+                        error!("[CLIENT {}] Process msg: {}",  self.address, e);
                         break
                     }
                 }
@@ -86,7 +86,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
         self.subscription.unsubscribe().await;
     }
 
-    pub async fn process_msg_from_p2p(&mut self, msg: &Privmsg) -> Result<()> {
+    pub async fn process_msg(&mut self, msg: &Privmsg) -> Result<()> {
         info!("[P2P] Received: {}", msg.to_string().trim());
 
         let mut msg = msg.clone();
@@ -213,16 +213,10 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
             }
 
             // Send dm messages in buffer
-            let privmsgs_buffer = self.buffers.privmsgs.lock().await;
-            for msg in privmsgs_buffer.iter() {
-                let is_dm = msg.target == self.irc_config.nickname ||
-                    (msg.nickname == self.irc_config.nickname && !msg.target.starts_with('#'));
-
-                if is_dm {
-                    self.notify_clients.notify_by_id(msg.clone(), self.subscription.get_id()).await;
-                }
+            let privmsgs = self.buffers.privmsgs.lock().await.clone();
+            for msg in privmsgs.iter() {
+                self.process_msg(msg).await?;
             }
-            drop(privmsgs_buffer);
         }
         Ok(())
     }
@@ -520,19 +514,7 @@ impl<C: AsyncRead + AsyncWrite + Send + Unpin + 'static> IrcClient<C> {
                 self.reply(&j).await?;
                 self.reply(&t).await?;
             }
-
-            // Send messages in buffer
-            if !self.irc_config.capabilities.get("no-history").unwrap() {
-                for msg in self.buffers.privmsgs.lock().await.iter() {
-                    if msg.target == *chan {
-                        self.notify_clients
-                            .notify_by_id(msg.clone(), self.subscription.get_id())
-                            .await;
-                    }
-                }
-            }
         }
-        self.on_receive_names(channels).await?;
         Ok(())
     }
 }

+ 0 - 11
src/system/subscriber.rs

@@ -77,17 +77,6 @@ impl<T: Clone> Subscriber<T> {
         }
     }
 
-    pub async fn notify_by_id(&self, message_result: T, id: u64) {
-        if let Some(sub) = (*self.subs.lock().await).get(&id) {
-            match sub.send(message_result.clone()).await {
-                Ok(()) => {}
-                Err(err) => {
-                    warn!("Error returned sending message in notify_by_id() call! {}", err);
-                }
-            }
-        }
-    }
-
     pub async fn notify_with_exclude(&self, message_result: T, exclude_list: &[SubscriptionId]) {
         for (id, sub) in (*self.subs.lock().await).iter() {
             if exclude_list.contains(id) {