seed_session.rs 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  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. let settings = {
  18. let p2p = self.p2p.upgrade().unwrap();
  19. p2p.settings()
  20. };
  21. if settings.skip_seed_sync {
  22. info!("Configured to skip seed synchronization process.");
  23. return Ok(());
  24. }
  25. // if cached addresses then quit
  26. // if seeds empty then seeding required but empty
  27. if settings.seeds.is_empty() {
  28. error!("Seeding is required but no seeds are configured.");
  29. return Err(NetError::OperationFailed);
  30. }
  31. let mut tasks = Vec::new();
  32. for seed in settings.seeds.clone() {
  33. tasks.push(executor.spawn(self.clone().start_seed(seed, executor.clone())));
  34. }
  35. for task in tasks {
  36. // Ignore errors
  37. let _ = task.await;
  38. }
  39. // Seed process complete
  40. // TODO: check increase count of address
  41. Ok(())
  42. }
  43. async fn start_seed(
  44. self: Arc<Self>,
  45. seed: SocketAddr,
  46. executor: Arc<Executor<'_>>,
  47. ) -> NetResult<()> {
  48. let (hosts, settings) = {
  49. let p2p = self.p2p.upgrade().unwrap();
  50. (p2p.hosts(), p2p.settings())
  51. };
  52. let connector = Connector::new(settings.clone());
  53. match connector.connect(seed).await {
  54. Ok(channel) => {
  55. // Blacklist goes here
  56. info!("Connected seed [{}]", seed);
  57. self.clone()
  58. .register_channel(channel.clone(), executor.clone())
  59. .await?;
  60. self.attach_protocols(channel, hosts, settings, executor)
  61. .await
  62. }
  63. Err(err) => {
  64. info!("Failure contacting seed [{}]: {}", seed, err);
  65. Err(err)
  66. }
  67. }
  68. }
  69. async fn attach_protocols(
  70. self: Arc<Self>,
  71. channel: ChannelPtr,
  72. hosts: HostsPtr,
  73. settings: SettingsPtr,
  74. executor: Arc<Executor<'_>>,
  75. ) -> NetResult<()> {
  76. let protocol_ping = ProtocolPing::new(channel.clone(), settings.clone());
  77. protocol_ping.start(executor.clone()).await;
  78. let protocol_seed = ProtocolSeed::new(channel.clone(), hosts, settings.clone());
  79. protocol_seed.start(executor.clone()).await?;
  80. channel.stop().await;
  81. Ok(())
  82. }
  83. }
  84. impl Session for SeedSession {
  85. fn p2p(&self) -> Arc<P2p> {
  86. self.p2p.upgrade().unwrap()
  87. }
  88. }