inbound_session.rs 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112
  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, ProtocolAddress, ProtocolSeed};
  7. use crate::net::sessions::Session;
  8. use crate::net::{Acceptor, AcceptorPtr};
  9. use crate::net::{ChannelPtr, Connector, HostsPtr, P2p, SettingsPtr};
  10. use crate::system::{StoppableTask, StoppableTaskPtr};
  11. pub struct InboundSession {
  12. p2p: Weak<P2p>,
  13. acceptor: AcceptorPtr,
  14. accept_task: StoppableTaskPtr,
  15. }
  16. impl InboundSession {
  17. pub fn new(p2p: Weak<P2p>) -> Arc<Self> {
  18. let settings = {
  19. let p2p = p2p.upgrade().unwrap();
  20. p2p.settings()
  21. };
  22. let acceptor = Acceptor::new(settings);
  23. Arc::new(Self { p2p, acceptor, accept_task: StoppableTask::new() })
  24. }
  25. pub fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> NetResult<()> {
  26. match self.p2p().settings().inbound {
  27. Some(accept_addr) => {
  28. self.clone().start_accept_session(accept_addr, executor.clone())?;
  29. }
  30. None => {
  31. info!("Not configured for accepting incoming connections.");
  32. return Ok(());
  33. }
  34. }
  35. self.accept_task.clone().start(
  36. self.clone().channel_sub_loop(executor.clone()),
  37. // Ignore stop handler
  38. |_| { async {} },
  39. NetError::ServiceStopped,
  40. executor);
  41. Ok(())
  42. }
  43. pub async fn stop(&self) {
  44. self.acceptor.stop().await;
  45. }
  46. fn start_accept_session(
  47. self: Arc<Self>,
  48. accept_addr: SocketAddr,
  49. executor: Arc<Executor<'_>>,
  50. ) -> NetResult<()> {
  51. info!("Starting inbound session on {}", accept_addr);
  52. let result = self.acceptor.clone().start(accept_addr, executor);
  53. if let Err(err) = result {
  54. error!("Error starting listener: {}", err);
  55. }
  56. result
  57. }
  58. async fn channel_sub_loop(self: Arc<Self>, executor: Arc<Executor<'_>>) -> NetResult<()> {
  59. let channel_sub = self.acceptor.clone().subscribe().await;
  60. loop {
  61. let channel = (*channel_sub.receive().await).clone()?;
  62. // Spawn a detached task to process the channel
  63. // This will just perform the channel setup then exit.
  64. executor.spawn(self.clone().setup_channel(channel, executor.clone())).detach();
  65. }
  66. }
  67. async fn setup_channel(self: Arc<Self>, channel: ChannelPtr, executor: Arc<Executor<'_>>) -> NetResult<()> {
  68. info!("Connected inbound [{}]", channel.address());
  69. self.clone()
  70. .register_channel(channel.clone(), executor.clone())
  71. .await?;
  72. let settings = self.p2p.upgrade().unwrap().settings();
  73. self.attach_protocols(channel, settings, executor)
  74. .await
  75. }
  76. async fn attach_protocols(
  77. self: Arc<Self>,
  78. channel: ChannelPtr,
  79. settings: SettingsPtr,
  80. executor: Arc<Executor<'_>>,
  81. ) -> NetResult<()> {
  82. let protocol_ping = ProtocolPing::new(channel.clone(), settings.clone());
  83. protocol_ping.start(executor.clone()).await;
  84. let protocol_addr = ProtocolAddress::new(channel, settings);
  85. protocol_addr.start(executor).await;
  86. Ok(())
  87. }
  88. }
  89. impl Session for InboundSession {
  90. fn p2p(&self) -> Arc<P2p> {
  91. self.p2p.upgrade().unwrap()
  92. }
  93. }