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

add protocol_address which does addr/get_addr hosts sync

narodnik 5 лет назад
Родитель
Сommit
961d586cef

+ 8 - 16
src/bin/dfi.rs

@@ -1,15 +1,15 @@
 #[macro_use]
 extern crate clap;
 use async_executor::Executor;
+use async_native_tls::TlsAcceptor;
 use async_std::sync::Mutex;
 use easy_parallel::Parallel;
+use http_types::{Request, Response, StatusCode};
 use serde_json::json;
+use smol::Async;
 use std::net::SocketAddr;
-use std::sync::Arc;
 use std::net::TcpListener;
-use async_native_tls::TlsAcceptor;
-use http_types::{Request, Response, StatusCode};
-use smol::Async;
+use std::sync::Arc;
 
 use sapvi::{net, Result};
 
@@ -248,10 +248,6 @@ async fn start2(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<
 
 struct ProgramOptions {
     network_settings: net::Settings,
-    accept_addr: Option<SocketAddr>,
-    seed_addrs: Vec<SocketAddr>,
-    manual_connects: Vec<SocketAddr>,
-    connection_slots: u32,
     log_path: Box<std::path::PathBuf>,
 }
 
@@ -306,19 +302,15 @@ impl ProgramOptions {
 
         Ok(ProgramOptions {
             network_settings: net::Settings {
-                inbound: accept_addr.clone(),
+                inbound: accept_addr,
                 outbound_connections: connection_slots,
                 connect_timeout_seconds: 10,
                 channel_handshake_seconds: 2,
                 channel_heartbeat_seconds: 10,
-                external_addr: accept_addr.clone(),
-                peers: manual_connects.clone(),
-                seeds: seed_addrs.clone(),
+                external_addr: accept_addr,
+                peers: manual_connects,
+                seeds: seed_addrs,
             },
-            accept_addr,
-            seed_addrs,
-            manual_connects,
-            connection_slots,
             log_path,
         })
     }

+ 1 - 1
src/net/connector.rs

@@ -1,5 +1,5 @@
 use futures::FutureExt;
-use smol::{Async};
+use smol::Async;
 use std::net::{SocketAddr, TcpStream};
 
 use crate::net::error::{NetError, NetResult};

+ 5 - 1
src/net/hosts.rs

@@ -24,11 +24,15 @@ impl Hosts {
         self.addrs.lock().await.extend(addrs)
     }
 
-    pub async fn load(&self) -> Option<SocketAddr> {
+    pub async fn load_single(&self) -> Option<SocketAddr> {
         self.addrs
             .lock()
             .await
             .choose(&mut rand::thread_rng())
             .cloned()
     }
+
+    pub async fn load_all(&self) -> Vec<SocketAddr> {
+        self.addrs.lock().await.clone()
+    }
 }

+ 0 - 2
src/net/mod.rs

@@ -11,7 +11,6 @@ pub mod hosts;
 pub mod messages;
 pub mod p2p;
 pub mod protocols;
-pub mod proxy;
 pub mod sessions;
 pub mod settings;
 pub mod utility;
@@ -24,5 +23,4 @@ pub use connector::Connector;
 pub use hosts::{Hosts, HostsPtr};
 pub use message_subscriber::{MessageSubscriber, MessageSubscription};
 pub use p2p::P2p;
-pub use proxy::Proxy;
 pub use settings::{Settings, SettingsPtr};

+ 1 - 3
src/net/p2p.rs

@@ -6,14 +6,13 @@ use std::sync::Arc;
 
 use crate::net::error::NetResult;
 use crate::net::sessions::{InboundSession, SeedSession};
-use crate::net::{Channel, ChannelPtr, Connector, Hosts, HostsPtr, Settings, SettingsPtr};
+use crate::net::{Channel, ChannelPtr, Hosts, HostsPtr, Settings, SettingsPtr};
 
 pub type Pending<T> = Mutex<HashMap<SocketAddr, Arc<T>>>;
 
 pub type P2pPtr = Arc<P2p>;
 
 pub struct P2p {
-    pending_connects: Pending<Connector>,
     pending_channels: Pending<Channel>,
     hosts: HostsPtr,
     settings: SettingsPtr,
@@ -23,7 +22,6 @@ impl P2p {
     pub fn new(settings: Settings) -> Arc<Self> {
         let settings = Arc::new(settings);
         Arc::new(Self {
-            pending_connects: Mutex::new(HashMap::new()),
             pending_channels: Mutex::new(HashMap::new()),
             hosts: Hosts::new(settings.clone()),
             settings,

+ 52 - 3
src/net/protocols/protocol_address.rs

@@ -1,24 +1,73 @@
 use smol::Executor;
 use std::sync::Arc;
 
+use crate::net::error::NetResult;
+use crate::net::messages;
 use crate::net::protocols::{ProtocolJobsManager, ProtocolJobsManagerPtr};
-use crate::net::{ChannelPtr, SettingsPtr};
+use crate::net::{ChannelPtr, HostsPtr, SettingsPtr};
 
 pub struct ProtocolAddress {
     channel: ChannelPtr,
+    hosts: HostsPtr,
     settings: SettingsPtr,
 
     jobsman: ProtocolJobsManagerPtr,
 }
 
 impl ProtocolAddress {
-    pub fn new(channel: ChannelPtr, settings: SettingsPtr) -> Arc<Self> {
+    pub fn new(channel: ChannelPtr, hosts: HostsPtr, settings: SettingsPtr) -> Arc<Self> {
         Arc::new(Self {
             channel: channel.clone(),
+            hosts,
             settings,
             jobsman: ProtocolJobsManager::new(channel),
         })
     }
 
-    pub async fn start(self: Arc<Self>, _executor: Arc<Executor<'_>>) {}
+    pub async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) {
+        self.jobsman.clone().start(executor.clone());
+        self.jobsman
+            .clone()
+            .spawn(self.clone().handle_receive_addrs(), executor.clone())
+            .await;
+        self.jobsman
+            .clone()
+            .spawn(self.clone().handle_receive_get_addrs(), executor)
+            .await;
+
+        // Send get_address message
+        let get_addrs = messages::Message::GetAddrs(messages::GetAddrsMessage {});
+        let _ = self.channel.clone().send(get_addrs).await;
+    }
+
+    async fn handle_receive_addrs(self: Arc<Self>) -> NetResult<()> {
+        let addrs_sub = self
+            .channel
+            .clone()
+            .subscribe_msg(messages::PacketType::Addrs)
+            .await;
+
+        loop {
+            let addrs_msg = receive_message!(addrs_sub, messages::Message::Addrs);
+
+            self.hosts.store(addrs_msg.addrs.clone()).await;
+        }
+    }
+
+    async fn handle_receive_get_addrs(self: Arc<Self>) -> NetResult<()> {
+        let get_addrs_sub = self
+            .channel
+            .clone()
+            .subscribe_msg(messages::PacketType::GetAddrs)
+            .await;
+
+        loop {
+            let _get_addrs = receive_message!(get_addrs_sub, messages::Message::GetAddrs);
+
+            let addrs = messages::Message::Addrs(messages::AddrsMessage {
+                addrs: self.hosts.load_all().await,
+            });
+            self.channel.clone().send(addrs).await?;
+        }
+    }
 }

+ 1 - 1
src/net/protocols/protocol_ping.rs

@@ -1,6 +1,6 @@
 use log::*;
 use rand::Rng;
-use smol::{Executor};
+use smol::Executor;
 use std::sync::Arc;
 
 use crate::net::error::{NetError, NetResult};

+ 0 - 12
src/net/proxy.rs

@@ -1,12 +0,0 @@
-use smol::{Async};
-use std::net::{TcpStream};
-
-pub struct Proxy {
-    stream: Async<TcpStream>,
-}
-
-impl Proxy {
-    pub fn new(stream: Async<TcpStream>) -> Self {
-        Self { stream }
-    }
-}

+ 6 - 6
src/net/sessions/inbound_session.rs

@@ -7,7 +7,7 @@ use crate::net::error::{NetError, NetResult};
 use crate::net::protocols::{ProtocolAddress, ProtocolPing};
 use crate::net::sessions::Session;
 use crate::net::{Acceptor, AcceptorPtr};
-use crate::net::{ChannelPtr, P2p, SettingsPtr};
+use crate::net::{ChannelPtr, P2p};
 use crate::system::{StoppableTask, StoppableTaskPtr};
 
 pub struct InboundSession {
@@ -95,21 +95,21 @@ impl InboundSession {
             .register_channel(channel.clone(), executor.clone())
             .await?;
 
-        let settings = self.p2p.upgrade().unwrap().settings();
-
-        self.attach_protocols(channel, settings, executor).await
+        self.attach_protocols(channel, executor).await
     }
 
     async fn attach_protocols(
         self: Arc<Self>,
         channel: ChannelPtr,
-        settings: SettingsPtr,
         executor: Arc<Executor<'_>>,
     ) -> NetResult<()> {
+        let settings = self.p2p().settings().clone();
+        let hosts = self.p2p().hosts().clone();
+
         let protocol_ping = ProtocolPing::new(channel.clone(), settings.clone());
         protocol_ping.start(executor.clone()).await;
 
-        let protocol_addr = ProtocolAddress::new(channel, settings);
+        let protocol_addr = ProtocolAddress::new(channel, hosts, settings);
         protocol_addr.start(executor).await;
 
         Ok(())