Kaynağa Gözat

net/channel: Close transports during shutdown

x 2 hafta önce
ebeveyn
işleme
bd26864618
2 değiştirilmiş dosya ile 66 ekleme ve 13 silme
  1. 37 2
      src/net/channel.rs
  2. 29 11
      src/net/tests.rs

+ 37 - 2
src/net/channel.rs

@@ -23,7 +23,7 @@ use std::{
         atomic::{AtomicBool, Ordering::SeqCst},
         Arc,
     },
-    time::UNIX_EPOCH,
+    time::{Duration, UNIX_EPOCH},
 };
 
 use darkfi_serial::{
@@ -54,7 +54,10 @@ use super::{
 };
 use crate::{
     net::BanPolicy,
-    system::{msleep, Publisher, PublisherPtr, StoppableTask, StoppableTaskPtr, Subscription},
+    system::{
+        msleep, timeout::timeout, Publisher, PublisherPtr, StoppableTask, StoppableTaskPtr,
+        Subscription,
+    },
     util::{logger::verbose, time::NanoTimestamp},
     Error, Result,
 };
@@ -62,6 +65,8 @@ use crate::{
 /// Atomic pointer to async channel
 pub type ChannelPtr = Arc<Channel>;
 
+const TRANSPORT_CLOSE_TIMEOUT: Duration = Duration::from_secs(5);
+
 /// Channel debug info
 #[derive(Clone, Debug, SerialEncodable, SerialDecodable)]
 pub struct ChannelInfo {
@@ -221,6 +226,34 @@ impl Channel {
         self.receive_task.stop_nowait();
     }
 
+    /// Gracefully close the underlying transport without allowing an
+    /// unresponsive peer to stall channel shutdown indefinitely.
+    async fn close_transport(&self) {
+        let close = async {
+            let writer = &mut *self.writer.lock().await;
+            writer.close().await
+        };
+
+        match timeout(TRANSPORT_CLOSE_TIMEOUT, close).await {
+            Ok(Ok(())) => {}
+            Ok(Err(err)) if err.kind() == io::ErrorKind::NotConnected => {}
+            Ok(Err(err)) => {
+                verbose!(
+                    target: "net::channel::close_transport",
+                    "[P2P] Failed closing channel transport {}: {err}",
+                    self.display_address(),
+                );
+            }
+            Err(_) => {
+                verbose!(
+                    target: "net::channel::close_transport",
+                    "[P2P] Timed out closing channel transport {}",
+                    self.display_address(),
+                );
+            }
+        }
+    }
+
     #[cfg(test)]
     pub(crate) async fn cleanup_task_count(&self) -> usize {
         self.cleanup_tasks.lock().await.len()
@@ -455,6 +488,8 @@ impl Channel {
             task.await;
         }
 
+        self.close_transport().await;
+
         self.p2p().untrack_channel(self.info.id);
 
         debug!(target: "net::channel::handle_stop", "[END] {self:?}");

+ 29 - 11
src/net/tests.rs

@@ -768,14 +768,24 @@ fn p2p_shutdown_drains_channels_across_restarts() {
 async fn p2p_shutdown_drains_channels_across_restarts_real(ex: Arc<Executor<'static>>) {
     const LIFECYCLE_COUNT: usize = 3;
 
+    for scheme in ["tcp", "tcp+tls"] {
+        p2p_shutdown_drains_transport_across_restarts(scheme, ex.clone(), LIFECYCLE_COUNT).await;
+    }
+}
+
+async fn p2p_shutdown_drains_transport_across_restarts(
+    scheme: &str,
+    ex: Arc<Executor<'static>>,
+    lifecycle_count: usize,
+) {
     let port = get_random_available_port();
-    let listen_url = Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap();
+    let listen_url = Url::parse(&format!("{scheme}://127.0.0.1:{port}")).unwrap();
     let server_settings = Settings {
         localnet: true,
         inbound_addrs: vec![listen_url.clone()],
         inbound_connections: 8,
         outbound_connections: 0,
-        active_profiles: vec!["tcp".to_string()],
+        active_profiles: vec![scheme.to_string()],
         ..Default::default()
     };
     let client_settings = Settings {
@@ -783,14 +793,14 @@ async fn p2p_shutdown_drains_channels_across_restarts_real(ex: Arc<Executor<'sta
         peers: vec![listen_url],
         inbound_connections: 0,
         outbound_connections: 0,
-        active_profiles: vec!["tcp".to_string()],
+        active_profiles: vec![scheme.to_string()],
         ..Default::default()
     };
 
     let server = P2p::new(server_settings, ex.clone()).await.unwrap();
     let client = P2p::new(client_settings, ex).await.unwrap();
 
-    for _ in 0..LIFECYCLE_COUNT {
+    for _ in 0..lifecycle_count {
         server.clone().start().await.unwrap();
         client.clone().start().await.unwrap();
 
@@ -802,20 +812,28 @@ async fn p2p_shutdown_drains_channels_across_restarts_real(ex: Arc<Executor<'sta
         .await
         .expect("manual connection was not established");
 
-        let channels = server
-            .hosts()
-            .channels()
-            .into_iter()
-            .chain(client.hosts().channels())
-            .collect::<Vec<_>>();
+        let server_channels = server.hosts().channels();
+        let client_channels = client.hosts().channels();
 
         client.stop().await;
+
+        // Stopping the client must close its transport even though client_channels
+        // keeps the Channel object alive. The server should observe EOF and stop
+        // its corresponding channel without waiting for its own P2P shutdown.
+        timeout(Duration::from_secs(5), async {
+            while server_channels.iter().any(|channel| !channel.is_stopped()) {
+                Timer::after(Duration::from_millis(10)).await;
+            }
+        })
+        .await
+        .expect("stopped client retained an open channel transport");
+
         server.stop().await;
 
         assert!(client.hosts().channels().is_empty());
         assert!(server.hosts().channels().is_empty());
         assert_eq!(server.session_inbound().connection_count().await, 0);
-        for channel in channels {
+        for channel in server_channels.into_iter().chain(client_channels) {
             assert!(channel.is_stopped());
             assert_eq!(channel.cleanup_task_count().await, 0);
         }