Ver Fonte

net/direct_session: add dnet support

epiphany há 9 meses atrás
pai
commit
0982c8ca72
4 ficheiros alterados com 111 adições e 13 exclusões
  1. 24 0
      src/net/dnet.rs
  2. 33 12
      src/net/session/direct_session.rs
  3. 14 1
      src/net/session/mod.rs
  4. 40 0
      src/rpc/from_impl.rs

+ 24 - 0
src/net/dnet.rs

@@ -83,6 +83,26 @@ pub struct OutboundPeerDiscovery {
     pub state: &'static str,
 }
 
+#[derive(Clone, Debug)]
+pub struct DirectConnecting {
+    pub connect_addr: Url,
+}
+
+#[derive(Clone, Debug)]
+pub struct DirectConnected {
+    pub connect_addr: Url,
+    pub addr: Url,
+    pub channel_id: u32,
+}
+
+#[derive(Clone, Debug)]
+pub struct DirectDisconnected {
+    pub connect_addr: Url,
+    pub err: String,
+}
+
+pub type DirectPeerDiscovery = OutboundPeerDiscovery;
+
 #[derive(Clone, Debug)]
 pub enum DnetEvent {
     SendMessage(MessageInfo),
@@ -94,4 +114,8 @@ pub enum DnetEvent {
     OutboundSlotConnected(OutboundSlotConnected),
     OutboundSlotDisconnected(OutboundSlotDisconnected),
     OutboundPeerDiscovery(OutboundPeerDiscovery),
+    DirectConnecting(DirectConnecting),
+    DirectConnected(DirectConnected),
+    DirectDisconnected(DirectDisconnected),
+    DirectPeerDiscovery(DirectPeerDiscovery),
 }

+ 33 - 12
src/net/session/direct_session.rs

@@ -43,19 +43,15 @@ use url::Url;
 use super::{
     super::{
         connector::Connector,
+        dnet::{self, dnetev, DnetEvent},
+        hosts::{HostColor, HostState},
+        message::GetAddrsMessage,
         p2p::{P2p, P2pPtr},
     },
     Session, SessionBitFlag, SESSION_DIRECT,
 };
 use crate::{
-    net::{
-        dnet,
-        dnet::{dnetev, DnetEvent},
-        hosts::HostState,
-        message::GetAddrsMessage,
-        session::HostColor,
-        ChannelPtr,
-    },
+    net::ChannelPtr,
     system::{sleep, timeout::timeout, CondVar, PublisherPtr, StoppableTask, StoppableTaskPtr},
     Error, Result,
 };
@@ -321,6 +317,10 @@ impl ChannelBuilder {
             return Err(e)
         }
 
+        dnetev!(self, DirectConnecting, {
+            connect_addr: addr.clone(),
+        });
+
         match self.connector().connect(addr).await {
             Ok((_, channel)) => {
                 info!(
@@ -329,6 +329,12 @@ impl ChannelBuilder {
                     channel.display_address()
                 );
 
+                dnetev!(self, DirectConnected, {
+                    connect_addr: channel.info.connect_addr.clone(),
+                    addr: channel.display_address().clone(),
+                    channel_id: channel.info.id
+                });
+
                 // Register the new channel
                 match self
                     .session()
@@ -346,6 +352,11 @@ impl ChannelBuilder {
                             channel.display_address(),
                         );
 
+                        dnetev!(self, DirectDisconnected, {
+                            connect_addr: channel.info.connect_addr.clone(),
+                            err: e.to_string()
+                        });
+
                         // Free up this addr for future operations.
                         if let Err(e) = self.session().p2p().hosts().unregister(channel.address()) {
                             warn!(target: "net::direct_session", "[P2P] Error while unregistering addr={}, err={e}", channel.display_address());
@@ -361,6 +372,11 @@ impl ChannelBuilder {
                     "[P2P] Unable to connect to direct outbound: {e}",
                 );
 
+                dnetev!(self, DirectDisconnected, {
+                    connect_addr: addr.clone(),
+                    err: e.to_string()
+                });
+
                 // Free up this addr for future operations.
                 if let Err(e) = self.session().p2p().hosts().unregister(addr) {
                     warn!(target: "net::direct_session", "[P2P] Error while unregistering addr={addr}, err={e}");
@@ -436,7 +452,7 @@ impl PeerDiscovery {
 
         let mut current_attempt = 0;
         loop {
-            dnetev!(self, OutboundPeerDiscovery, {
+            dnetev!(self, DirectPeerDiscovery, {
                 attempt: current_attempt,
                 state: "wait",
             });
@@ -460,7 +476,7 @@ impl PeerDiscovery {
                     "[P2P] [PEER DISCOVERY] Sleeping and trying again. Attempt {current_attempt}"
                 );
 
-                dnetev!(self, OutboundPeerDiscovery, {
+                dnetev!(self, DirectPeerDiscovery, {
                     attempt: current_attempt,
                     state: "sleep",
                 });
@@ -474,6 +490,11 @@ impl PeerDiscovery {
             // whitelist, or greylist.
             let mut channel = None;
             if !self.p2p().is_connected() {
+                dnetev!(self, DirectPeerDiscovery, {
+                    attempt: current_attempt,
+                    state: "newchan",
+                });
+
                 for color in [HostColor::Gold, HostColor::White, HostColor::Grey].iter() {
                     if let Some((entry, _)) = self
                         .p2p()
@@ -496,7 +517,7 @@ impl PeerDiscovery {
                     target: "net::direct_session::peer_discovery()",
                     "[P2P] [PEER DISCOVERY] Asking peers for new peers to connect to...");
 
-                dnetev!(self, OutboundPeerDiscovery, {
+                dnetev!(self, DirectPeerDiscovery, {
                     attempt: current_attempt,
                     state: "getaddr",
                 });
@@ -548,7 +569,7 @@ impl PeerDiscovery {
                     target: "net::direct_session::peer_discovery()",
                     "[P2P] [PEER DISCOVERY] Asking seeds for new peers to connect to...");
 
-                dnetev!(self, OutboundPeerDiscovery, {
+                dnetev!(self, DirectPeerDiscovery, {
                     attempt: current_attempt,
                     state: "seed",
                 });

+ 14 - 1
src/net/session/mod.rs

@@ -25,7 +25,13 @@ use async_trait::async_trait;
 use log::{debug, error, trace};
 use smol::Executor;
 
-use super::{channel::ChannelPtr, hosts::HostColor, p2p::P2pPtr, protocol::ProtocolVersion};
+use super::{
+    channel::ChannelPtr,
+    dnet::{self, dnetev, DnetEvent},
+    hosts::HostColor,
+    p2p::P2pPtr,
+    protocol::ProtocolVersion,
+};
 use crate::{system::Subscription, Error, Result};
 
 pub mod inbound_session;
@@ -110,6 +116,13 @@ pub async fn remove_sub_on_stop(
         }
     }
 
+    if type_id & SESSION_DIRECT != 0 {
+        dnetev!(p2p.session_direct(), DirectDisconnected, {
+            connect_addr: channel.info.connect_addr.clone(),
+            err: "Channel stopped".to_string()
+        });
+    }
+
     if !p2p.is_connected() {
         hosts.disconnect_publisher.notify(Error::NetworkNotConnected).await;
     }

+ 40 - 0
src/rpc/from_impl.rs

@@ -95,6 +95,34 @@ impl From<net::dnet::OutboundPeerDiscovery> for JsonValue {
     }
 }
 
+#[cfg(feature = "net")]
+impl From<net::dnet::DirectConnecting> for JsonValue {
+    fn from(info: net::dnet::DirectConnecting) -> JsonValue {
+        json_map([("connect_addr", JsonStr(info.connect_addr.to_string()))])
+    }
+}
+
+#[cfg(feature = "net")]
+impl From<net::dnet::DirectConnected> for JsonValue {
+    fn from(info: net::dnet::DirectConnected) -> JsonValue {
+        json_map([
+            ("connect_addr", JsonStr(info.connect_addr.to_string())),
+            ("addr", JsonStr(info.addr.to_string())),
+            ("channel_id", JsonNum(info.channel_id.into())),
+        ])
+    }
+}
+
+#[cfg(feature = "net")]
+impl From<net::dnet::DirectDisconnected> for JsonValue {
+    fn from(info: net::dnet::DirectDisconnected) -> JsonValue {
+        json_map([
+            ("connect_addr", JsonStr(info.connect_addr.to_string())),
+            ("err", JsonStr(info.err)),
+        ])
+    }
+}
+
 #[cfg(feature = "net")]
 impl From<net::dnet::DnetEvent> for JsonValue {
     fn from(event: net::dnet::DnetEvent) -> JsonValue {
@@ -126,6 +154,18 @@ impl From<net::dnet::DnetEvent> for JsonValue {
             net::dnet::DnetEvent::OutboundPeerDiscovery(info) => {
                 json_map([("event", json_str("outbound_peer_discovery")), ("info", info.into())])
             }
+            net::dnet::DnetEvent::DirectConnecting(info) => {
+                json_map([("event", json_str("direct_connecting")), ("info", info.into())])
+            }
+            net::dnet::DnetEvent::DirectConnected(info) => {
+                json_map([("event", json_str("direct_connected")), ("info", info.into())])
+            }
+            net::dnet::DnetEvent::DirectDisconnected(info) => {
+                json_map([("event", json_str("direct_disconnected")), ("info", info.into())])
+            }
+            net::dnet::DnetEvent::DirectPeerDiscovery(info) => {
+                json_map([("event", json_str("direct_peer_discovery")), ("info", info.into())])
+            }
         }
     }
 }