outbound_session.rs 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2024 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. //! Outbound connections session. Manages the creation of outbound sessions.
  19. //! Used to create an outbound session and to stop and start the session.
  20. //!
  21. //! Class consists of a weak pointer to the p2p interface and a vector of
  22. //! outbound connection slots. Using a weak pointer to p2p allows us to
  23. //! avoid circular dependencies. The vector of slots is wrapped in a mutex
  24. //! lock. This is switched on every time we instantiate a connection slot
  25. //! and insures that no other part of the program uses the slots at the
  26. //! same time.
  27. use std::{
  28. sync::{
  29. atomic::{AtomicU32, Ordering},
  30. Arc, Weak,
  31. },
  32. time::{Duration, Instant},
  33. };
  34. use async_trait::async_trait;
  35. use log::{debug, error, info, warn};
  36. use smol::lock::Mutex;
  37. use url::Url;
  38. use super::{
  39. super::{
  40. channel::ChannelPtr,
  41. connector::Connector,
  42. dnet::{self, dnetev, DnetEvent},
  43. hosts::store::{HostColor, HostState},
  44. message::GetAddrsMessage,
  45. p2p::{P2p, P2pPtr},
  46. },
  47. Session, SessionBitFlag, SESSION_OUTBOUND,
  48. };
  49. use crate::{
  50. system::{sleep, timeout::timeout, CondVar, LazyWeak, StoppableTask, StoppableTaskPtr},
  51. Error, Result,
  52. };
  53. pub type OutboundSessionPtr = Arc<OutboundSession>;
  54. /// Defines outbound connections session.
  55. pub struct OutboundSession {
  56. /// Weak pointer to parent p2p object
  57. pub(in crate::net) p2p: LazyWeak<P2p>,
  58. /// Outbound connection slots
  59. slots: Mutex<Vec<Arc<Slot>>>,
  60. /// Peer discovery task
  61. peer_discovery: Arc<PeerDiscovery>,
  62. }
  63. impl OutboundSession {
  64. /// Create a new outbound session.
  65. pub(crate) fn new() -> OutboundSessionPtr {
  66. let self_ = Arc::new(Self {
  67. p2p: LazyWeak::new(),
  68. slots: Mutex::new(Vec::new()),
  69. peer_discovery: PeerDiscovery::new(),
  70. });
  71. self_.peer_discovery.session.init(self_.clone());
  72. self_
  73. }
  74. /// Start the outbound session. Runs the channel connect loop.
  75. pub(crate) async fn start(self: Arc<Self>) {
  76. let n_slots = self.p2p().settings().outbound_connections;
  77. info!(target: "net::outbound_session", "[P2P] Starting {} outbound connection slots.", n_slots);
  78. // Activate mutex lock on connection slots.
  79. let mut slots = self.slots.lock().await;
  80. let self_ = Arc::downgrade(&self);
  81. for i in 0..n_slots as u32 {
  82. let slot = Slot::new(self_.clone(), i);
  83. slot.clone().start().await;
  84. slots.push(slot);
  85. }
  86. self.peer_discovery.clone().start().await;
  87. }
  88. /// Stops the outbound session.
  89. pub(crate) async fn stop(&self) {
  90. let slots = &*self.slots.lock().await;
  91. for slot in slots {
  92. slot.clone().stop().await;
  93. }
  94. self.peer_discovery.clone().stop().await;
  95. }
  96. pub async fn slot_info(&self) -> Vec<u32> {
  97. let mut info = Vec::new();
  98. let slots = &*self.slots.lock().await;
  99. for slot in slots {
  100. info.push(slot.channel_id.load(Ordering::Relaxed));
  101. }
  102. info
  103. }
  104. fn wakeup_peer_discovery(&self) {
  105. self.peer_discovery.notify()
  106. }
  107. async fn wakeup_slots(&self) {
  108. let slots = &*self.slots.lock().await;
  109. for slot in slots {
  110. slot.notify();
  111. }
  112. }
  113. }
  114. #[async_trait]
  115. impl Session for OutboundSession {
  116. fn p2p(&self) -> P2pPtr {
  117. self.p2p.upgrade()
  118. }
  119. fn type_id(&self) -> SessionBitFlag {
  120. SESSION_OUTBOUND
  121. }
  122. }
  123. pub struct Slot {
  124. slot: u32,
  125. process: StoppableTaskPtr,
  126. wakeup_self: CondVar,
  127. session: Weak<OutboundSession>,
  128. // For debugging
  129. channel_id: AtomicU32,
  130. }
  131. impl Slot {
  132. fn new(session: Weak<OutboundSession>, slot: u32) -> Arc<Self> {
  133. Arc::new(Self {
  134. slot,
  135. process: StoppableTask::new(),
  136. wakeup_self: CondVar::new(),
  137. session,
  138. channel_id: AtomicU32::new(0),
  139. })
  140. }
  141. async fn start(self: Arc<Self>) {
  142. // TODO: way too many clones, look into making this implicit. See implicit-clone crate
  143. let ex = self.p2p().executor();
  144. self.process.clone().start(
  145. async move {
  146. self.run().await;
  147. unreachable!();
  148. },
  149. // Ignore stop handler
  150. |_| async {},
  151. Error::NetworkServiceStopped,
  152. ex,
  153. );
  154. }
  155. async fn stop(self: Arc<Self>) {
  156. self.process.stop().await
  157. }
  158. /// Address selection algorithm that works as follows: up to
  159. /// anchor_count, select from the anchorlist. Up to white_count,
  160. /// select from the whitelist. For all other slots, select from
  161. /// the greylist.
  162. /// If we didn't find an address with this selection logic, downgrade
  163. /// our preferences. Up to anchor_count, select from the whitelist,
  164. /// up until white_count, select from the greylist.
  165. /// If we still didn't find an address, select from the greylist. In
  166. /// all other cases, return an empty vector. This will trigger
  167. /// fetch_addrs() to return None and initiate peer discovery.
  168. /* NOTE: Selecting from the greylist for some % of the slots is
  169. necessary and healthy since we require the network retains some
  170. unreliable connections. A network that purely favors uptime over
  171. unreliable connections may be vulnerable to sybil by attackers with
  172. good uptime.*/
  173. async fn fetch_addrs_with_preference(&self, preference: usize) -> Vec<(Url, u64)> {
  174. let slot = self.slot;
  175. let settings = self.p2p().settings();
  176. let hosts = &self.p2p().hosts().container;
  177. let white_count = settings.white_connect_count;
  178. let anchor_count = settings.anchor_connect_count;
  179. let transports = &settings.allowed_transports;
  180. let transport_mixing = settings.transport_mixing;
  181. debug!(target: "net::outbound_session::fetch_addrs_with_preference()",
  182. "slot={}, preference={}", slot, preference);
  183. match preference {
  184. // Highest preference that corresponds to the anchor and white count preference set in
  185. // Settings.
  186. 0 => {
  187. if slot < anchor_count {
  188. hosts.fetch(HostColor::Gold, transports, transport_mixing).await
  189. } else if slot < white_count {
  190. hosts.fetch(HostColor::White, transports, transport_mixing).await
  191. } else {
  192. hosts.fetch(HostColor::Grey, transports, transport_mixing).await
  193. }
  194. }
  195. // Reduced preference in case we don't have sufficient hosts to satisfy our highest
  196. // preference.
  197. 1 => {
  198. if slot < anchor_count {
  199. hosts.fetch(HostColor::White, transports, transport_mixing).await
  200. } else if slot < white_count {
  201. hosts.fetch(HostColor::Grey, transports, transport_mixing).await
  202. } else {
  203. vec![]
  204. }
  205. }
  206. // Lowest preference if we still haven't been able to find a host.
  207. 2 => {
  208. if slot < anchor_count {
  209. hosts.fetch(HostColor::Grey, transports, transport_mixing).await
  210. } else {
  211. vec![]
  212. }
  213. }
  214. _ => {
  215. panic!()
  216. }
  217. }
  218. }
  219. // Fetch an address we can connect to acccording to the white and anchor connection counts
  220. // configured in Settings.
  221. async fn fetch_addrs(&self) -> Option<(Url, u64)> {
  222. let hosts = self.p2p().hosts();
  223. // First select an addresses that match our white and anchor requirements configured in
  224. // Settings.
  225. let preference = 0;
  226. let addrs = self.fetch_addrs_with_preference(preference).await;
  227. if !addrs.is_empty() {
  228. return hosts.check_addrs(addrs).await;
  229. }
  230. // If no addresses were returned, go for the second best thing (white and grey).
  231. let preference = 1;
  232. let addrs = self.fetch_addrs_with_preference(preference).await;
  233. if !addrs.is_empty() {
  234. return hosts.check_addrs(addrs).await;
  235. }
  236. // If we still have no addresses, go for the least favored option.
  237. let preference = 2;
  238. let addrs = self.fetch_addrs_with_preference(preference).await;
  239. if !addrs.is_empty() {
  240. return hosts.check_addrs(addrs).await;
  241. }
  242. // If we still don't have an address, return None and do peer discovery.
  243. None
  244. }
  245. // We first try to make connections to the addresses on our anchor list. We then find some
  246. // whitelist connections according to the whitelist percent default. Finally, any remaining
  247. // connections we make from the greylist.
  248. async fn run(self: Arc<Self>) {
  249. let hosts = self.p2p().hosts();
  250. loop {
  251. // Activate the slot
  252. debug!(
  253. target: "net::outbound_session::try_connect()",
  254. "[P2P] Finding a host to connect to for outbound slot #{}",
  255. self.slot,
  256. );
  257. // Do peer discovery if we don't have a hostlist (first time connecting
  258. // to the network).
  259. if hosts.container.is_empty(HostColor::Grey).await {
  260. dnetev!(self, OutboundSlotSleeping, {
  261. slot: self.slot,
  262. });
  263. self.wakeup_self.reset();
  264. // Peer discovery
  265. self.session().wakeup_peer_discovery();
  266. // Wait to be woken up by peer discovery
  267. self.wakeup_self.wait().await;
  268. continue
  269. }
  270. let addr = if let Some(addr) = self.fetch_addrs().await {
  271. debug!(target: "net::outbound_session::run()", "Fetched address: {:?}", addr);
  272. addr
  273. } else {
  274. debug!(target: "net::outbound_session::run()", "No address found! Activating peer discovery...");
  275. dnetev!(self, OutboundSlotSleeping, {
  276. slot: self.slot,
  277. });
  278. self.wakeup_self.reset();
  279. // Peer discovery
  280. self.session().wakeup_peer_discovery();
  281. // Wait to be woken up by peer discovery
  282. self.wakeup_self.wait().await;
  283. continue
  284. };
  285. let host = addr.0;
  286. let last_seen = addr.1;
  287. let slot = self.slot;
  288. info!(
  289. target: "net::outbound_session::try_connect()",
  290. "[P2P] Connecting outbound slot #{} [{}]",
  291. slot, host,
  292. );
  293. dnetev!(self, OutboundSlotConnecting, {
  294. slot,
  295. addr: host.clone(),
  296. });
  297. let (addr, channel) = match self.try_connect(host.clone(), last_seen).await {
  298. Ok(connect_info) => connect_info,
  299. Err(err) => {
  300. debug!(
  301. target: "net::outbound_session::try_connect()",
  302. "[P2P] Outbound slot #{} connection failed: {}",
  303. slot, err
  304. );
  305. dnetev!(self, OutboundSlotDisconnected, {
  306. slot,
  307. err: err.to_string()
  308. });
  309. self.channel_id.store(0, Ordering::Relaxed);
  310. continue
  311. }
  312. };
  313. info!(
  314. target: "net::outbound_session::try_connect()",
  315. "[P2P] Outbound slot #{} connected [{}]",
  316. slot, addr
  317. );
  318. dnetev!(self, OutboundSlotConnected, {
  319. slot: self.slot,
  320. addr: addr.clone(),
  321. channel_id: channel.info.id
  322. });
  323. // At this point we've managed to connect.
  324. let stop_sub = channel.subscribe_stop().await.expect("Channel should not be stopped");
  325. // Setup new channel
  326. if let Err(err) =
  327. self.session().register_channel(channel.clone(), self.p2p().executor()).await
  328. {
  329. info!(
  330. target: "net::outbound_session",
  331. "[P2P] Outbound slot #{} disconnected: {}",
  332. slot, err
  333. );
  334. dnetev!(self, OutboundSlotDisconnected, {
  335. slot: self.slot,
  336. err: err.to_string()
  337. });
  338. self.channel_id.store(0, Ordering::Relaxed);
  339. continue
  340. }
  341. self.channel_id.store(channel.info.id, Ordering::Relaxed);
  342. // Add this connection to the anchorlist
  343. hosts
  344. .move_host(&addr, last_seen, HostColor::Gold, Some(channel.clone()))
  345. .await
  346. .unwrap();
  347. // Wait for channel to close
  348. stop_sub.receive().await;
  349. self.channel_id.store(0, Ordering::Relaxed);
  350. }
  351. }
  352. /// Start making an outbound connection, using provided [`Connector`].
  353. /// Tries to find a valid address to connect to, otherwise does peer
  354. /// discovery. The peer discovery loops until some peer we can connect
  355. /// to is found. Once connected, registers the channel, removes it from
  356. /// the list of pending channels, and starts sending messages across the
  357. /// channel. In case of any failures, a network error is returned and the
  358. /// main connect loop (parent of this function) will iterate again.
  359. async fn try_connect(&self, addr: Url, last_seen: u64) -> Result<(Url, ChannelPtr)> {
  360. let parent = Arc::downgrade(&self.session());
  361. let connector = Connector::new(self.p2p().settings(), parent);
  362. match connector.connect(&addr).await {
  363. Ok((addr_final, channel)) => Ok((addr_final, channel)),
  364. Err(e) => {
  365. debug!(
  366. target: "net::outbound_session::try_connect()",
  367. "[P2P] Unable to connect outbound slot #{} [{}]: {}",
  368. self.slot, addr, e
  369. );
  370. // At this point we failed to connect. We'll downgrade this peer now.
  371. self.p2p().hosts().move_host(&addr, last_seen, HostColor::Grey, None).await?;
  372. // Mark its state as Suspend, which sends it to the Refinery for processing.
  373. self.p2p().hosts().try_register(addr.clone(), HostState::Suspend).await.unwrap();
  374. // Notify that channel processing failed
  375. self.p2p().hosts().channel_subscriber.notify(Err(Error::ConnectFailed)).await;
  376. Err(Error::ConnectFailed)
  377. }
  378. }
  379. }
  380. fn notify(&self) {
  381. self.wakeup_self.notify()
  382. }
  383. fn session(&self) -> OutboundSessionPtr {
  384. self.session.upgrade().unwrap()
  385. }
  386. fn p2p(&self) -> P2pPtr {
  387. self.session().p2p()
  388. }
  389. }
  390. /// Defines a common interface for multiple peer discovery processes.
  391. /* NOTE: Currently only one Peer Discovery implementation exists. Making
  392. Peer Discovery generic enables us to support network swarming, since
  393. the peer discovery process will differ depending on whether it occurs
  394. on the overlay network or a subnet.*/
  395. #[async_trait]
  396. pub trait PeerDiscoveryBase {
  397. async fn start(self: Arc<Self>);
  398. async fn stop(self: Arc<Self>);
  399. async fn run(self: Arc<Self>);
  400. async fn wait(&self) -> bool;
  401. fn notify(&self);
  402. fn session(&self) -> OutboundSessionPtr;
  403. fn p2p(&self) -> P2pPtr;
  404. }
  405. /// Main PeerDiscovery process that loops through connected channels
  406. /// and sends out a `GetAddrs` when it is active.
  407. struct PeerDiscovery {
  408. process: StoppableTaskPtr,
  409. wakeup_self: CondVar,
  410. session: LazyWeak<OutboundSession>,
  411. }
  412. impl PeerDiscovery {
  413. fn new() -> Arc<Self> {
  414. Arc::new(Self {
  415. process: StoppableTask::new(),
  416. wakeup_self: CondVar::new(),
  417. session: LazyWeak::new(),
  418. })
  419. }
  420. }
  421. #[async_trait]
  422. impl PeerDiscoveryBase for PeerDiscovery {
  423. async fn start(self: Arc<Self>) {
  424. let ex = self.p2p().executor();
  425. self.process.clone().start(
  426. async move {
  427. self.run().await;
  428. unreachable!();
  429. },
  430. // Ignore stop handler
  431. |_| async {},
  432. Error::NetworkServiceStopped,
  433. ex,
  434. );
  435. }
  436. async fn stop(self: Arc<Self>) {
  437. self.process.stop().await
  438. }
  439. /// Activate peer discovery if not active already. This will loop through all
  440. /// connected P2P channels and send out a `GetAddrs` message to request more
  441. /// peers. Other parts of the P2P stack will then handle the incoming addresses
  442. /// and place them in the hosts list.
  443. /// This function will also sleep `Settings::outbound_peer_discovery_attempt_time` seconds
  444. /// after broadcasting in order to let the P2P stack receive and work through
  445. /// the addresses it is expecting.
  446. async fn run(self: Arc<Self>) {
  447. let mut current_attempt = 0;
  448. loop {
  449. dnetev!(self, OutboundPeerDiscovery, {
  450. attempt: current_attempt,
  451. state: "wait",
  452. });
  453. // wait to be woken up by notify()
  454. let sleep_was_instant = self.wait().await;
  455. let p2p = self.p2p();
  456. if sleep_was_instant {
  457. // Try again
  458. current_attempt += 1;
  459. } else {
  460. // reset back to start
  461. current_attempt = 1;
  462. }
  463. if current_attempt >= 4 {
  464. debug!("current attempt: {}", current_attempt);
  465. info!(
  466. target: "net::outbound_session::peer_discovery()",
  467. "[P2P] Sleeping and trying again..."
  468. );
  469. dnetev!(self, OutboundPeerDiscovery, {
  470. attempt: current_attempt,
  471. state: "sleep",
  472. });
  473. sleep(p2p.settings().outbound_peer_discovery_cooloff_time).await;
  474. current_attempt = 1;
  475. }
  476. // First 2 times try sending GetAddr to the network.
  477. // 3rd time do a seed sync.
  478. if p2p.is_connected().await && current_attempt <= 2 {
  479. // Broadcast the GetAddrs message to all active channels.
  480. // If we have no active channels, we will perform a SeedSyncSession instead.
  481. info!(
  482. target: "net::outbound_session::peer_discovery()",
  483. "[P2P] Requesting addrs from active channels. Attempt: {}",
  484. current_attempt
  485. );
  486. dnetev!(self, OutboundPeerDiscovery, {
  487. attempt: current_attempt,
  488. state: "getaddr",
  489. });
  490. let get_addrs = GetAddrsMessage {
  491. max: p2p.settings().outbound_connections as u32,
  492. transports: p2p.settings().allowed_transports.clone(),
  493. };
  494. p2p.broadcast(&get_addrs).await;
  495. // Wait for a hosts store update event
  496. let store_sub = self.p2p().hosts().subscribe_store().await;
  497. let result = timeout(
  498. Duration::from_secs(p2p.settings().outbound_peer_discovery_attempt_time),
  499. store_sub.receive(),
  500. )
  501. .await;
  502. match result {
  503. Ok(addrs_len) => {
  504. info!(
  505. target: "net::outbound_session::peer_discovery()",
  506. "[P2P] Discovered {} addrs", addrs_len
  507. );
  508. }
  509. Err(_) => {
  510. warn!(
  511. target: "net::outbound_session::peer_discovery()",
  512. "[P2P] Peer discovery waiting for addrs timed out."
  513. );
  514. // TODO: Just do seed next time
  515. }
  516. }
  517. // TODO: check every subscribe() call has a corresponding unsubscribe()
  518. store_sub.unsubscribe().await;
  519. } else {
  520. info!(
  521. target: "net::outbound_session::peer_discovery()",
  522. "[P2P] Seeding hosts. Attempt: {}",
  523. current_attempt
  524. );
  525. dnetev!(self, OutboundPeerDiscovery, {
  526. attempt: current_attempt,
  527. state: "seed",
  528. });
  529. match p2p.clone().seed().await {
  530. Ok(()) => {
  531. info!(
  532. target: "net::outbound_session::peer_discovery()",
  533. "[P2P] Seeding hosts successful."
  534. );
  535. }
  536. Err(err) => {
  537. error!(
  538. target: "net::outbound_session::peer_discovery()",
  539. "[P2P] Network reseed failed: {}", err,
  540. );
  541. }
  542. }
  543. }
  544. self.wakeup_self.reset();
  545. self.session().wakeup_slots().await;
  546. // Give some time for new connections to be established
  547. sleep(p2p.settings().outbound_peer_discovery_attempt_time).await;
  548. }
  549. }
  550. async fn wait(&self) -> bool {
  551. let wakeup_start = Instant::now();
  552. self.wakeup_self.wait().await;
  553. let wakeup_end = Instant::now();
  554. let epsilon = Duration::from_millis(200);
  555. wakeup_end - wakeup_start <= epsilon
  556. }
  557. fn notify(&self) {
  558. self.wakeup_self.notify()
  559. }
  560. fn session(&self) -> OutboundSessionPtr {
  561. self.session.upgrade()
  562. }
  563. fn p2p(&self) -> P2pPtr {
  564. self.session().p2p()
  565. }
  566. }