seed_session.rs 3.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103
  1. use async_executor::Executor;
  2. use log::*;
  3. use std::net::SocketAddr;
  4. use std::sync::{Arc, Weak};
  5. use crate::error::{Error, Result};
  6. use crate::net::sessions::Session;
  7. use crate::net::{ChannelPtr, HostsPtr, Connector, P2p, SettingsPtr};
  8. use crate::net::protocols::{ProtocolPing, ProtocolSeed};
  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<'_>>) -> Result<()> {
  17. let settings = {
  18. let p2p = self.p2p.upgrade().unwrap();
  19. p2p.settings()
  20. };
  21. // if cached addresses then quit
  22. // if seeds empty then seeding required but empty
  23. if settings.seeds.is_empty() {
  24. error!("Seeding is required but no seeds are configured.");
  25. return Err(Error::OperationFailed);
  26. }
  27. let mut tasks = Vec::new();
  28. for seed in settings.seeds.clone() {
  29. tasks.push(executor.spawn(self.clone().start_seed(seed, executor.clone())));
  30. }
  31. for task in tasks {
  32. // Ignore errors
  33. let _ = task.await;
  34. }
  35. // Seed process complete
  36. // TODO: check increase count of address
  37. Ok(())
  38. }
  39. async fn start_seed(self: Arc<Self>, seed: SocketAddr, executor: Arc<Executor<'_>>) -> Result<()> {
  40. let (hosts, settings) = {
  41. let p2p = self.p2p.upgrade().unwrap();
  42. (p2p.hosts(), p2p.settings())
  43. };
  44. let connector = Connector::new(settings.clone());
  45. match connector.connect(seed).await {
  46. Ok(channel) => {
  47. // Blacklist goes here
  48. info!("Connected seed [{}]", seed);
  49. self.clone().register_channel(channel.clone(), executor.clone()).await?;
  50. self.attach_protocols(channel, hosts, settings, executor).await
  51. }
  52. Err(err) => {
  53. info!("Failure contacting seed [{}]: {}", seed, err);
  54. Err(err)
  55. }
  56. }
  57. }
  58. async fn register_channel(self: Arc<Self>, channel: ChannelPtr, executor: Arc<Executor<'_>>) -> Result<()> {
  59. let handshake_task = self.perform_handshake_protocols(channel.clone(), executor.clone());
  60. // start channel
  61. channel.start(executor);
  62. handshake_task.await
  63. }
  64. async fn attach_protocols(self: Arc<Self>, channel: ChannelPtr, hosts: HostsPtr, settings: SettingsPtr, executor: Arc<Executor<'_>>) -> Result<()> {
  65. let protocol_ping = ProtocolPing::new(channel.clone(), settings.clone());
  66. let ping_task = protocol_ping.start(executor.clone());
  67. let protocol_seed = ProtocolSeed::new(channel, hosts, settings.clone());
  68. protocol_seed.start(executor.clone()).await?;
  69. // Close the ping task now we finished.
  70. // TODO: channel drop should trigger this automatically anyway via the stop signal
  71. ping_task.cancel().await;
  72. Ok(())
  73. }
  74. }
  75. impl Session for SeedSession {
  76. fn p2p(&self) -> Arc<P2p> {
  77. self.p2p.upgrade().unwrap()
  78. }
  79. }