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

abstract privmsg stuff into a proper Protocol object that uses the jobs manager to ensure the running tasks are properly closed when the channel disconnects.

narodnik 4 лет назад
Родитель
Сommit
90d859e424
2 измененных файлов с 68 добавлено и 35 удалено
  1. 67 35
      src/bin/ircd.rs
  2. 1 0
      src/net/mod.rs

+ 67 - 35
src/bin/ircd.rs

@@ -168,41 +168,6 @@ impl Decodable for PrivMsg {
     }
 }
 
-async fn channel_loop(
-    p2p: net::P2pPtr,
-    sender: async_channel::Sender<Arc<PrivMsg>>,
-    executor: Arc<Executor<'_>>,
-) -> Result<()> {
-    debug!("CHANNEL SUBS LOOP");
-    let new_channel_sub = p2p.subscribe_channel().await;
-
-    loop {
-        let channel = new_channel_sub.receive().await?;
-
-        debug!("NEWCHANNEL");
-
-        let message_subsytem = channel.get_message_subsystem();
-        message_subsytem.add_dispatch::<PrivMsg>().await;
-
-        debug!("ADDED DISPATCH");
-
-        let privmsg_sub = channel.subscribe_msg::<PrivMsg>().await?;
-        executor.spawn(catch_privmsgs(privmsg_sub, sender.clone())).detach();
-    }
-}
-
-async fn catch_privmsgs(
-    privmsg_sub: net::MessageSubscription<PrivMsg>,
-    sender: async_channel::Sender<Arc<PrivMsg>>,
-) -> Result<()> {
-    loop {
-        let privmsg = privmsg_sub.receive().await?;
-        debug!("GOTIT {:?}", privmsg);
-        sender.send(privmsg).await.expect("send message");
-        debug!("SENT OVER THE TUBES");
-    }
-}
-
 async fn process(
     recvr: async_channel::Receiver<Arc<PrivMsg>>,
     stream: Async<TcpStream>,
@@ -271,6 +236,73 @@ async fn process_user_input(
     }
 }
 
+struct ProtocolPrivMsg {
+    notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
+    privmsg_sub: net::MessageSubscription<PrivMsg>,
+    jobsman: net::ProtocolJobsManagerPtr,
+}
+
+impl ProtocolPrivMsg {
+    async fn new(
+        channel: net::ChannelPtr,
+        notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
+    ) -> Arc<Self> {
+        let message_subsytem = channel.get_message_subsystem();
+        message_subsytem.add_dispatch::<PrivMsg>().await;
+
+        debug!("ADDED DISPATCH");
+
+        let privmsg_sub =
+            channel.subscribe_msg::<PrivMsg>().await.expect("Missing PrivMsg dispatcher!");
+
+        Arc::new(Self {
+            notify_queue_sender,
+            privmsg_sub,
+            jobsman: net::ProtocolJobsManager::new("PrivMsgProtocol", channel),
+        })
+    }
+
+    async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) {
+        debug!(target: "ircd", "ProtocolPrivMsg::start() [START]");
+        self.jobsman.clone().start(executor.clone());
+        self.jobsman.clone().spawn(self.clone().handle_receive_privmsg(), executor.clone()).await;
+        debug!(target: "ircd", "ProtocolPrivMsg::start() [END]");
+    }
+
+    async fn handle_receive_privmsg(self: Arc<Self>) -> Result<()> {
+        debug!(target: "ircd", "ProtocolAddress::handle_receive_privmsg() [START]");
+        loop {
+            let privmsg = self.privmsg_sub.receive().await?;
+
+            debug!(
+                target: "ircd",
+                "ProtocolPrivMsg::handle_receive_privmsg() received {:?}",
+                privmsg
+            );
+
+            self.notify_queue_sender.send(privmsg).await.expect("notify_queue_sender send failed!");
+        }
+    }
+}
+
+async fn channel_loop(
+    p2p: net::P2pPtr,
+    sender: async_channel::Sender<Arc<PrivMsg>>,
+    executor: Arc<Executor<'_>>,
+) -> Result<()> {
+    debug!("CHANNEL SUBS LOOP");
+    let new_channel_sub = p2p.subscribe_channel().await;
+
+    loop {
+        let channel = new_channel_sub.receive().await?;
+
+        debug!("NEWCHANNEL");
+
+        let protocol_privmsg = ProtocolPrivMsg::new(channel, sender.clone()).await;
+        protocol_privmsg.start(executor.clone()).await;
+    }
+}
+
 async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
     let listener = match Async::<TcpListener>::bind(options.irc_accept_addr) {
         Ok(listener) => listener,

+ 1 - 0
src/net/mod.rs

@@ -95,4 +95,5 @@ pub use hosts::{Hosts, HostsPtr};
 pub use message_subscriber::MessageSubscription;
 pub use messages::Message;
 pub use p2p::{P2p, P2pPtr};
+pub use protocols::{ProtocolJobsManager, ProtocolJobsManagerPtr};
 pub use settings::{Settings, SettingsPtr};