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

Switch IRCd to using net protocol registry

narodnik 4 лет назад
Родитель
Сommit
3b914bfe36
2 измененных файлов с 43 добавлено и 34 удалено
  1. 27 26
      bin/ircd/src/main.rs
  2. 16 8
      bin/ircd/src/protocol_privmsg.rs

+ 27 - 26
bin/ircd/src/main.rs

@@ -98,25 +98,6 @@ async fn process_user_input(
     Ok(())
 }
 
-async fn channel_loop(
-    p2p: net::P2pPtr,
-    sender: async_channel::Sender<Arc<PrivMsg>>,
-    seen_privmsg_ids: SeenPrivMsgIdsPtr,
-    executor: Arc<Executor<'_>>,
-) -> Result<()> {
-    let new_channel_sub = p2p.subscribe_channel().await;
-
-    loop {
-        let channel = new_channel_sub.receive().await?;
-
-        let protocol_privmsg =
-            ProtocolPrivMsg::new(channel, sender.clone(), seen_privmsg_ids.clone(), p2p.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,
@@ -145,7 +126,28 @@ async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<(
 
     let seen_privmsg_ids = SeenPrivMsgIds::new();
 
+    //
+    // PrivMsg protocol
+    //
     let p2p = net::P2p::new(options.network_settings).await;
+    let registry = p2p.protocol_registry();
+
+    let (sender, recvr) = async_channel::unbounded();
+    let seen_privmsg_ids2 = seen_privmsg_ids.clone();
+    let sender2 = sender.clone();
+    registry.register(
+        !net::SESSION_SEED,
+        move |channel, p2p| {
+            let sender = sender2.clone();
+            let seen_privmsg_ids = seen_privmsg_ids2.clone();
+            async move {
+                ProtocolPrivMsg::new(channel, sender, seen_privmsg_ids, p2p).await
+            }
+        }).await;
+
+    //
+    // p2p network main instance
+    //
     // Performs seed session
     p2p.clone().start(executor.clone()).await?;
     // Actual main p2p session
@@ -159,13 +161,9 @@ async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<(
         })
         .detach();
 
-    let (sender, recvr) = async_channel::unbounded();
-    // for now the p2p and channel sub sessions just run forever
-    // so detach them as background processes.
-    executor
-        .spawn(channel_loop(p2p.clone(), sender, seen_privmsg_ids.clone(), executor.clone()))
-        .detach();
-
+    //
+    // RPC interface
+    //
     let ex2 = executor.clone();
     let ex3 = ex2.clone();
     let rpc_interface = Arc::new(JsonRpcInterface {});
@@ -173,6 +171,9 @@ async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<(
         .spawn(async move { listen_and_serve(server_config, rpc_interface, ex3).await })
         .detach();
 
+    //
+    // IRC instance
+    //
     loop {
         let (stream, peer_addr) = match listener.accept().await {
             Ok((s, a)) => (s, a),

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

@@ -1,3 +1,4 @@
+use async_trait::async_trait;
 use async_executor::Executor;
 
 use darkfi::{net, Result};
@@ -20,7 +21,7 @@ impl ProtocolPrivMsg {
         notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
         seen_privmsg_ids: SeenPrivMsgIdsPtr,
         p2p: net::P2pPtr,
-    ) -> Arc<Self> {
+    ) -> net::ProtocolBasePtr {
         let message_subsytem = channel.get_message_subsystem();
         message_subsytem.add_dispatch::<PrivMsg>().await;
 
@@ -38,13 +39,6 @@ impl ProtocolPrivMsg {
         })
     }
 
-    pub 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 {
@@ -71,3 +65,17 @@ impl ProtocolPrivMsg {
         }
     }
 }
+
+#[async_trait]
+impl net::ProtocolBase for ProtocolPrivMsg {
+    /// Starts ping-pong keep-alive messages exchange. Runs ping-pong in the
+    /// protocol task manager, then queues the reply. Sends out a ping and
+    /// waits for pong reply. Waits for ping and replies with a pong.
+    async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
+        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]");
+        Ok(())
+    }
+}