seed_session.rs 4.6 KB

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