Selaa lähdekoodia

outbound_session: start() and stop() slots concurrently

use FuturesUnordered to poll slot.start() and slot.stop() in a
non-blocking manner.  This should be more efficient than waiting for
each stop(), start() future to complete.
draoi 2 vuotta sitten
vanhempi
sitoutus
39350f8740
1 muutettua tiedostoa jossa 15 lisäystä ja 2 poistoa
  1. 15 2
      src/net/session/outbound_session.rs

+ 15 - 2
src/net/session/outbound_session.rs

@@ -35,6 +35,7 @@ use std::{
 };
 
 use async_trait::async_trait;
+use futures::stream::{FuturesUnordered, StreamExt};
 use log::{debug, error, info, warn};
 use smol::lock::Mutex;
 use url::Url;
@@ -86,14 +87,18 @@ impl OutboundSession {
         // Activate mutex lock on connection slots.
         let mut slots = self.slots.lock().await;
 
+        let mut futures = FuturesUnordered::new();
+
         let self_ = Arc::downgrade(&self);
 
         for i in 0..n_slots as u32 {
             let slot = Slot::new(self_.clone(), i);
-            slot.clone().start().await;
+            futures.push(slot.clone().start());
             slots.push(slot);
         }
 
+        while (futures.next().await).is_some() {}
+
         self.peer_discovery.clone().start().await;
     }
 
@@ -102,10 +107,18 @@ impl OutboundSession {
         debug!(target: "net::outbound_session", "Stopping outbound session");
         let slots = &*self.slots.lock().await;
 
+        let mut futures = FuturesUnordered::new();
+
         for slot in slots {
-            slot.clone().stop().await;
+            futures.push(slot.clone().stop());
         }
 
+        while (futures.next().await).is_some() {}
+
+        // TODO/ FIXME: Shutting down the slots triggers a seed sync in peer discovery
+        // (see outbound_session.rs:633).
+        // We should implement an Atomic Bool called stopped()
+        // (see channel.rs).
         self.peer_discovery.clone().stop().await;
     }