seed_session.rs 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139
  1. use async_executor::Executor;
  2. use futures::FutureExt;
  3. use log::*;
  4. use std::net::SocketAddr;
  5. use std::sync::{Arc, Weak};
  6. use crate::net::error::{NetError, NetResult};
  7. use crate::net::protocols::{ProtocolPing, ProtocolSeed};
  8. use crate::net::sessions::Session;
  9. use crate::net::utility::sleep;
  10. use crate::net::{ChannelPtr, Connector, HostsPtr, P2p, SettingsPtr};
  11. /// Seed connections session.
  12. pub struct SeedSession {
  13. p2p: Weak<P2p>,
  14. }
  15. impl SeedSession {
  16. /// Create a new seed session instance.
  17. pub fn new(p2p: Weak<P2p>) -> Arc<Self> {
  18. Arc::new(Self { p2p })
  19. }
  20. /// Start the seed session. Creates a new task for every seed connection and
  21. /// starts the seed on each task.
  22. pub async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> NetResult<()> {
  23. debug!(target: "net", "SeedSession::start() [START]");
  24. let settings = self.p2p().settings();
  25. if settings.seeds.is_empty() {
  26. warn!("Skipping seed sync process since no seeds are configured.");
  27. return Ok(());
  28. }
  29. // if cached addresses then quit
  30. let mut tasks = Vec::new();
  31. for (i, seed) in settings.seeds.iter().enumerate() {
  32. tasks.push(executor.spawn(self.clone().start_seed(i, seed.clone(), executor.clone())));
  33. }
  34. // This line loops through all the tasks and waits for them to finish.
  35. // But if the seed_query_timeout_seconds times out before they are finished,
  36. // then it will simply quit and the tasks will get dropped.
  37. futures::select! {
  38. _ = async move {
  39. for (i, task) in tasks.into_iter().enumerate() {
  40. // Ignore errors
  41. match task.await {
  42. Ok(()) => info!("Successfully queried seed #{}", i),
  43. Err(err) => warn!("Seed query #{} failed for reason: {}", i, err),
  44. }
  45. }
  46. }.fuse() => {
  47. }
  48. _ = sleep(settings.seed_query_timeout_seconds).fuse() => {
  49. error!("Querying seeds timed out");
  50. return Err(NetError::OperationFailed);
  51. }
  52. }
  53. // Seed process complete
  54. if self.p2p().hosts().is_empty().await {
  55. error!("Hosts pool still empty after seeding");
  56. return Err(NetError::OperationFailed);
  57. }
  58. debug!(target: "net", "SeedSession::start() [END]");
  59. Ok(())
  60. }
  61. /// Connects to a seed socket address. Registers a new channel with a network
  62. /// handshake, then starts the keep-alive messages and seed protocol.
  63. async fn start_seed(
  64. self: Arc<Self>,
  65. seed_index: usize,
  66. seed: SocketAddr,
  67. executor: Arc<Executor<'_>>,
  68. ) -> NetResult<()> {
  69. debug!(target: "net", "SeedSession::start_seed(i={}) [START]", seed_index);
  70. let (hosts, settings) = {
  71. let p2p = self.p2p.upgrade().unwrap();
  72. (p2p.hosts(), p2p.settings())
  73. };
  74. let connector = Connector::new(settings.clone());
  75. match connector.connect(seed).await {
  76. Ok(channel) => {
  77. // Blacklist goes here
  78. info!("Connected seed #{} [{}]", seed_index, seed);
  79. self.clone()
  80. .register_channel(channel.clone(), executor.clone())
  81. .await?;
  82. self.attach_protocols(channel, hosts, settings, executor)
  83. .await?;
  84. debug!(target: "net", "SeedSession::start_seed(i={}) [END]", seed_index);
  85. Ok(())
  86. }
  87. Err(err) => {
  88. info!(
  89. "Failure contacting seed #{} [{}]: {}",
  90. seed_index, seed, err
  91. );
  92. Err(err)
  93. }
  94. }
  95. }
  96. /// Starts keep-alive messages and seed protocol.
  97. async fn attach_protocols(
  98. self: Arc<Self>,
  99. channel: ChannelPtr,
  100. hosts: HostsPtr,
  101. settings: SettingsPtr,
  102. executor: Arc<Executor<'_>>,
  103. ) -> NetResult<()> {
  104. let protocol_ping = ProtocolPing::new(channel.clone(), settings.clone());
  105. protocol_ping.start(executor.clone()).await;
  106. let protocol_seed = ProtocolSeed::new(channel.clone(), hosts, settings.clone());
  107. // This will block until seed process is complete
  108. protocol_seed.start(executor.clone()).await?;
  109. channel.stop().await;
  110. Ok(())
  111. }
  112. }
  113. impl Session for SeedSession {
  114. fn p2p(&self) -> Arc<P2p> {
  115. self.p2p.upgrade().unwrap()
  116. }
  117. }