Kaynağa Gözat

darkirc: use StoppableTask instead of detach()

aggstam 3 yıl önce
ebeveyn
işleme
6ff0ddbd1d
2 değiştirilmiş dosya ile 107 ekleme ve 38 silme
  1. 38 17
      bin/darkirc/src/irc/server.rs
  2. 69 21
      bin/darkirc/src/main.rs

+ 38 - 17
bin/darkirc/src/irc/server.rs

@@ -33,7 +33,7 @@ use darkfi::{
         view::ViewPtr,
         view::ViewPtr,
     },
     },
     net::P2pPtr,
     net::P2pPtr,
-    system::SubscriberPtr,
+    system::{StoppableTask, SubscriberPtr},
     util::{path::expand_path, time::Timestamp},
     util::{path::expand_path, time::Timestamp},
     Error, Result,
     Error, Result,
 };
 };
@@ -85,31 +85,45 @@ impl IrcServer {
         let (msg_notifier, msg_recv) = smol::channel::unbounded();
         let (msg_notifier, msg_recv) = smol::channel::unbounded();
 
 
         // Listen to msgs from clients
         // Listen to msgs from clients
-        executor
-            .clone()
-            .spawn(Self::listen_to_msgs(
+        StoppableTask::new().start(
+            Self::listen_to_msgs(
                 self.p2p.clone(),
                 self.p2p.clone(),
                 self.model.clone(),
                 self.model.clone(),
                 self.seen.clone(),
                 self.seen.clone(),
                 msg_recv,
                 msg_recv,
                 self.missed_events.clone(),
                 self.missed_events.clone(),
                 self.clients_subscriptions.clone(),
                 self.clients_subscriptions.clone(),
-            ))
-            .detach();
+            ),
+            |res| async {
+                match res {
+                    Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                    Err(e) => error!(target: "darkirc::irc::server::start", "Failed starting listen to msgs: {}", e),
+                }
+            },
+            Error::DetachedTaskStopped,
+            executor.clone(),
+        );
 
 
         // Listen to msgs from View
         // Listen to msgs from View
-        executor
-            .clone()
-            .spawn(Self::listen_to_view(
+        StoppableTask::new().start(
+            Self::listen_to_view(
                 self.view.clone(),
                 self.view.clone(),
                 self.seen.clone(),
                 self.seen.clone(),
                 self.missed_events.clone(),
                 self.missed_events.clone(),
                 self.clients_subscriptions.clone(),
                 self.clients_subscriptions.clone(),
-            ))
-            .detach();
+            ),
+            |res| async {
+                match res {
+                    Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                    Err(e) => error!(target: "darkirc::irc::server::start", "Failed starting listen to view: {}", e),
+                }
+            },
+            Error::DetachedTaskStopped,
+            executor.clone(),
+        );
 
 
         // Start listening for new connections
         // Start listening for new connections
-        self.listen(msg_notifier, executor.clone()).await?;
+        self.listen(msg_notifier, executor).await?;
 
 
         Ok(())
         Ok(())
     }
     }
@@ -267,11 +281,18 @@ impl IrcServer {
         );
         );
 
 
         // Start listening and detach
         // Start listening and detach
-        executor
-            .spawn(async move {
-                client.listen().await;
-            })
-            .detach();
+        StoppableTask::new().start(
+            // Weird hack to prevent lifetimes hell
+            async move {client.listen().await; Ok(())},
+            |res| async {
+                match res {
+                    Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                    Err(e) => error!(target: "darkirc::irc::server::process_connection", "Failed starting client listen: {}", e),
+                }
+            },
+            Error::DetachedTaskStopped,
+            executor,
+        );
 
 
         Ok(())
         Ok(())
     }
     }

+ 69 - 21
bin/darkirc/src/main.rs

@@ -26,7 +26,7 @@ use async_std::{
 
 
 use chrono::{Duration, Utc};
 use chrono::{Duration, Utc};
 use irc::ClientSubMsg;
 use irc::ClientSubMsg;
-use log::{debug, info};
+use log::{debug, error, info};
 use rand::rngs::OsRng;
 use rand::rngs::OsRng;
 use structopt_toml::StructOptToml;
 use structopt_toml::StructOptToml;
 use tinyjson::JsonValue;
 use tinyjson::JsonValue;
@@ -41,9 +41,9 @@ use darkfi::{
     },
     },
     net,
     net,
     rpc::{jsonrpc::JsonSubscriber, server::listen_and_serve},
     rpc::{jsonrpc::JsonSubscriber, server::listen_and_serve},
-    system::{Subscriber, SubscriberPtr},
+    system::{StoppableTask, Subscriber, SubscriberPtr},
     util::{async_util::sleep, file::save_json_file, path::expand_path, time::Timestamp},
     util::{async_util::sleep, file::save_json_file, path::expand_path, time::Timestamp},
-    Result,
+    Error, Result,
 };
 };
 
 
 pub mod crypto;
 pub mod crypto;
@@ -227,33 +227,47 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
     ////////////////////
     ////////////////////
     // RPC interface setup
     // RPC interface setup
     ////////////////////
     ////////////////////
-
     let rpc_listen_addr = settings.rpc_listen.clone();
     let rpc_listen_addr = settings.rpc_listen.clone();
+    info!(target: "darkirc", "Starting JSON-RPC server on {}", rpc_listen_addr);
     let rpc_interface = Arc::new(JsonRpcInterface {
     let rpc_interface = Arc::new(JsonRpcInterface {
         addr: rpc_listen_addr.clone(),
         addr: rpc_listen_addr.clone(),
         p2p: p2p.clone(),
         p2p: p2p.clone(),
         dnet_sub: json_sub,
         dnet_sub: json_sub,
     });
     });
-
-    let _ex = executor.clone();
-    executor
-        .spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface, _ex).await })
-        .detach();
+    let rpc_task = StoppableTask::new();
+    rpc_task.clone().start(
+        listen_and_serve(rpc_listen_addr, rpc_interface, executor.clone()),
+        |res| async {
+            match res {
+                Ok(()) | Err(Error::RPCServerStopped) => { /* Do nothing */ }
+                Err(e) => error!(target: "darkirc", "Failed starting JSON-RPC server: {}", e),
+            }
+        },
+        Error::RPCServerStopped,
+        executor.clone(),
+    );
 
 
     ////////////////////
     ////////////////////
     // Start P2P network
     // Start P2P network
     ////////////////////
     ////////////////////
+    info!(target: "darkirc", "Starting P2P network");
     p2p.clone().start(executor.clone()).await?;
     p2p.clone().start(executor.clone()).await?;
-
-    // Run
-    let executor_cloned = executor.clone();
-    executor_cloned.spawn(p2p.clone().run(executor.clone())).detach();
+    StoppableTask::new().start(
+        p2p.clone().run(executor.clone()),
+        |res| async {
+            match res {
+                Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }
+                Err(e) => error!(target: "darkirc", "Failed starting P2P network: {}", e),
+            }
+        },
+        Error::P2PNetworkStopped,
+        executor.clone(),
+    );
 
 
     ////////////////////
     ////////////////////
     // IRC server
     // IRC server
     ////////////////////
     ////////////////////
-
-    // New irc server
+    info!(target: "darkirc", "Starting IRC server");
     let irc_server = IrcServer::new(
     let irc_server = IrcServer::new(
         settings.clone(),
         settings.clone(),
         p2p.clone(),
         p2p.clone(),
@@ -262,13 +276,37 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
         client_sub,
         client_sub,
     )
     )
     .await?;
     .await?;
-
-    // Start the irc server and detach it
-    let executor_cloned = executor.clone();
-    executor.spawn(async move { irc_server.start(executor_cloned).await }).detach();
+    let irc_server_task = StoppableTask::new();
+    let executor_ = executor.clone();
+    irc_server_task.clone().start(
+        // Weird hack to prevent lifetimes hell
+        async move { irc_server.start(executor_).await },
+        |res| async {
+            match res {
+                Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                Err(e) => error!(target: "darkirc", "Failed starting IRC server: {}", e),
+            }
+        },
+        Error::DetachedTaskStopped,
+        executor.clone(),
+    );
 
 
     // Reset root task
     // Reset root task
-    executor.spawn(async move { remove_old_events(model_clone2).await }).detach();
+    info!(target: "darkirc", "Starting remove old events task");
+    let remove_old_events_task = StoppableTask::new();
+    remove_old_events_task.clone().start(
+        remove_old_events(model_clone2),
+        |res| async {
+            match res {
+                Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
+                Err(e) => {
+                    error!(target: "darkirc", "Failed starting remove old events task: {}", e)
+                }
+            }
+        },
+        Error::DetachedTaskStopped,
+        executor,
+    );
 
 
     // Wait for termination signal
     // Wait for termination signal
     signals_handler.wait_termination(signals_task).await?;
     signals_handler.wait_termination(signals_task).await?;
@@ -276,7 +314,17 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
 
 
     model_clone.lock().await.save_tree(&datastore_path)?;
     model_clone.lock().await.save_tree(&datastore_path)?;
 
 
-    // stop p2p
+    info!(target: "darkirc", "Stopping JSON-RPC server...");
+    rpc_task.stop().await;
+
+    info!(target: "darkirc", "Stopping P2P network");
     p2p.stop().await;
     p2p.stop().await;
+
+    info!(target: "darkirc", "Stopping IRC server...");
+    irc_server_task.stop().await;
+
+    info!(target: "darkirc", "Stopping remove old events task...");
+    remove_old_events_task.stop().await;
+
     Ok(())
     Ok(())
 }
 }