seed_session.rs 3.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. use async_executor::Executor;
  2. use log::*;
  3. use std::net::SocketAddr;
  4. use std::sync::{Arc, Weak};
  5. use crate::net::error::{NetError, NetResult};
  6. use crate::net::protocols::{ProtocolPing, ProtocolSeed};
  7. use crate::net::sessions::Session;
  8. use crate::net::{ChannelPtr, Connector, HostsPtr, P2p, SettingsPtr};
  9. pub struct SeedSession {
  10. p2p: Weak<P2p>,
  11. }
  12. impl SeedSession {
  13. pub fn new(p2p: Weak<P2p>) -> Arc<Self> {
  14. Arc::new(Self { p2p })
  15. }
  16. pub async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> NetResult<()> {
  17. debug!(target: "net", "SeedSession::start() [START]");
  18. let settings = {
  19. let p2p = self.p2p.upgrade().unwrap();
  20. p2p.settings()
  21. };
  22. if settings.skip_seed_sync {
  23. info!("Configured to skip seed synchronization process.");
  24. return Ok(());
  25. }
  26. // if cached addresses then quit
  27. // if seeds empty then seeding required but empty
  28. if settings.seeds.is_empty() {
  29. error!("Seeding is required but no seeds are configured.");
  30. return Err(NetError::OperationFailed);
  31. }
  32. let mut tasks = Vec::new();
  33. for (i, seed) in settings.seeds.iter().enumerate() {
  34. tasks.push(executor.spawn(self.clone().start_seed(i, seed.clone(), executor.clone())));
  35. }
  36. for (i, task) in tasks.into_iter().enumerate() {
  37. // Ignore errors
  38. match task.await {
  39. Ok(()) => info!("Successfully queried seed #{}", i),
  40. Err(err) => warn!("Seed query #{} failed for reason: {}", i, err),
  41. }
  42. }
  43. // Seed process complete
  44. // TODO: check increase count of address
  45. debug!(target: "net", "SeedSession::start() [END]");
  46. Ok(())
  47. }
  48. async fn start_seed(
  49. self: Arc<Self>,
  50. seed_index: usize,
  51. seed: SocketAddr,
  52. executor: Arc<Executor<'_>>,
  53. ) -> NetResult<()> {
  54. debug!(target: "net", "SeedSession::start_seed(i={}) [START]", seed_index);
  55. let (hosts, settings) = {
  56. let p2p = self.p2p.upgrade().unwrap();
  57. (p2p.hosts(), p2p.settings())
  58. };
  59. let connector = Connector::new(settings.clone());
  60. match connector.connect(seed).await {
  61. Ok(channel) => {
  62. // Blacklist goes here
  63. info!("Connected seed #{} [{}]", seed_index, seed);
  64. self.clone()
  65. .register_channel(channel.clone(), executor.clone())
  66. .await?;
  67. self.attach_protocols(channel, hosts, settings, executor)
  68. .await?;
  69. debug!(target: "net", "SeedSession::start_seed(i={}) [END]", seed_index);
  70. Ok(())
  71. }
  72. Err(err) => {
  73. info!(
  74. "Failure contacting seed #{} [{}]: {}",
  75. seed_index, seed, err
  76. );
  77. Err(err)
  78. }
  79. }
  80. }
  81. async fn attach_protocols(
  82. self: Arc<Self>,
  83. channel: ChannelPtr,
  84. hosts: HostsPtr,
  85. settings: SettingsPtr,
  86. executor: Arc<Executor<'_>>,
  87. ) -> NetResult<()> {
  88. let protocol_ping = ProtocolPing::new(channel.clone(), settings.clone());
  89. protocol_ping.start(executor.clone()).await;
  90. let protocol_seed = ProtocolSeed::new(channel.clone(), hosts, settings.clone());
  91. // This will block until seed process is complete
  92. protocol_seed.start(executor.clone()).await?;
  93. channel.stop().await;
  94. Ok(())
  95. }
  96. }
  97. impl Session for SeedSession {
  98. fn p2p(&self) -> Arc<P2p> {
  99. self.p2p.upgrade().unwrap()
  100. }
  101. }