manual_session.rs 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 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. //! Manual connections session. Manages the creation of manual sessions.
  19. //! Used to create a manual session and to stop and start the session.
  20. //!
  21. //! A manual session is a type of outbound session in which we attempt
  22. //! connection to a predefined set of peers.
  23. //!
  24. //! Class consists of a weak pointer to the p2p interface and a vector of
  25. //! outbound connection slots. Using a weak pointer to p2p allows us to
  26. //! avoid circular dependencies. The vector of slots is wrapped in a mutex
  27. //! lock. This is switched on every time we instantiate a connection slot
  28. //! and insures that no other part of the program uses the slots at the
  29. //! same time.
  30. use async_std::sync::{Arc, Mutex, Weak};
  31. use async_trait::async_trait;
  32. use log::{info, warn};
  33. use smol::Executor;
  34. use url::Url;
  35. use super::{
  36. super::{
  37. channel::ChannelPtr,
  38. connector::Connector,
  39. p2p::{DnetInfo, P2p, P2pPtr},
  40. },
  41. Session, SessionBitFlag, SESSION_MANUAL,
  42. };
  43. use crate::{
  44. system::{StoppableTask, StoppableTaskPtr, Subscriber, SubscriberPtr},
  45. util::async_util::sleep,
  46. Error, Result,
  47. };
  48. pub type ManualSessionPtr = Arc<ManualSession>;
  49. /// Defines manual connections session.
  50. pub struct ManualSession {
  51. p2p: Weak<P2p>,
  52. connect_slots: Mutex<Vec<StoppableTaskPtr>>,
  53. /// Subscriber used to signal channels processing
  54. channel_subscriber: SubscriberPtr<Result<ChannelPtr>>,
  55. /// Flag to toggle channel_subscriber notifications
  56. notify: Mutex<bool>,
  57. }
  58. impl ManualSession {
  59. /// Create a new manual session.
  60. pub fn new(p2p: Weak<P2p>) -> ManualSessionPtr {
  61. Arc::new(Self {
  62. p2p,
  63. connect_slots: Mutex::new(vec![]),
  64. channel_subscriber: Subscriber::new(),
  65. notify: Mutex::new(false),
  66. })
  67. }
  68. /// Stops the manual session.
  69. pub async fn stop(&self) {
  70. let connect_slots = &*self.connect_slots.lock().await;
  71. for slot in connect_slots {
  72. slot.stop().await;
  73. }
  74. }
  75. /// Connect the manual session to the given address
  76. pub async fn connect(self: Arc<Self>, addr: Url, ex: Arc<Executor<'_>>) {
  77. let task = StoppableTask::new();
  78. task.clone().start(
  79. self.clone().channel_connect_loop(addr, ex.clone()),
  80. // Ignore stop handler
  81. |_| async {},
  82. Error::NetworkServiceStopped,
  83. ex,
  84. );
  85. self.connect_slots.lock().await.push(task);
  86. }
  87. /// Creates a connector object and tries to connect using it
  88. pub async fn channel_connect_loop(
  89. self: Arc<Self>,
  90. addr: Url,
  91. ex: Arc<Executor<'_>>,
  92. ) -> Result<()> {
  93. let parent = Arc::downgrade(&self);
  94. let settings = self.p2p().settings();
  95. let connector = Connector::new(settings.clone(), Arc::new(parent));
  96. let attempts = settings.manual_attempt_limit;
  97. let mut remaining = attempts;
  98. // Add the peer to list of pending channels
  99. self.p2p().add_pending(&addr).await;
  100. // Loop forever if attempts==0, otherwise loop attempts number of times.
  101. let mut tried_attempts = 0;
  102. loop {
  103. tried_attempts += 1;
  104. info!(
  105. target: "net::manual_session",
  106. "[P2P] Connecting to manual outbound [{}] (attempt #{})",
  107. addr, tried_attempts,
  108. );
  109. match connector.connect(addr.clone()).await {
  110. Ok(channel) => {
  111. info!(
  112. target: "net::manual_session",
  113. "[P2P] Manual outbound connected [{}]", addr,
  114. );
  115. let stop_sub =
  116. channel.subscribe_stop().await.expect("Channel should not be stopped");
  117. // Register the new channel
  118. self.register_channel(channel.clone(), ex.clone()).await?;
  119. // Channel is now connected but not yet setup
  120. // Remove pending lock since register_channel will add the channel to p2p
  121. self.p2p().remove_pending(&addr).await;
  122. // Notify that channel processing has finished
  123. if *self.notify.lock().await {
  124. self.channel_subscriber.notify(Ok(channel)).await;
  125. }
  126. // Wait for channel to close
  127. stop_sub.receive().await;
  128. info!(
  129. target: "net::manual_session",
  130. "[P2P] Manual outbound disconnected [{}]", addr,
  131. );
  132. // DEV NOTE: Here we can choose to attempt reconnection again
  133. return Ok(())
  134. }
  135. Err(e) => {
  136. warn!(
  137. target: "net::manual_session",
  138. "[P2P] Unable to connect to manual outbound [{}]: {}",
  139. addr, e,
  140. );
  141. }
  142. }
  143. // Wait and try again.
  144. // TODO: Should we notify about the failure now, or after all attempts
  145. // have failed?
  146. if *self.notify.lock().await {
  147. self.channel_subscriber.notify(Err(Error::ConnectFailed)).await;
  148. }
  149. remaining = if attempts == 0 { 1 } else { remaining - 1 };
  150. if remaining == 0 {
  151. break
  152. }
  153. info!(
  154. target: "net::manual_session",
  155. "[P2P] Waiting {} seconds until next manual outbound connection attempt [{}]",
  156. settings.outbound_connect_timeout, addr,
  157. );
  158. sleep(settings.outbound_connect_timeout).await;
  159. }
  160. warn!(
  161. target: "net::manual_session",
  162. "[P2P] Suspending manual connection to {} after {} failed attempts",
  163. addr, attempts,
  164. );
  165. self.p2p().remove_pending(&addr).await;
  166. Ok(())
  167. }
  168. /// Enable channel_subscriber notifications.
  169. pub async fn enable_notify(self: Arc<Self>) {
  170. *self.notify.lock().await = true;
  171. }
  172. /// Disable channel_subscriber notifications.
  173. pub async fn disable_notify(self: Arc<Self>) {
  174. *self.notify.lock().await = false;
  175. }
  176. }
  177. #[async_trait]
  178. impl Session for ManualSession {
  179. fn p2p(&self) -> P2pPtr {
  180. self.p2p.upgrade().unwrap()
  181. }
  182. fn type_id(&self) -> SessionBitFlag {
  183. SESSION_MANUAL
  184. }
  185. async fn dnet_info(&self) -> DnetInfo {
  186. todo!()
  187. }
  188. }