connector.rs 3.0 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495
  1. use async_std::sync::Arc;
  2. use std::{env, time::Duration};
  3. use log::error;
  4. use url::Url;
  5. use crate::{Error, Result};
  6. use super::{
  7. Channel, ChannelPtr, SessionWeakPtr, SettingsPtr, TcpTransport, TorTransport, Transport,
  8. TransportName,
  9. };
  10. /// Create outbound socket connections.
  11. pub struct Connector {
  12. settings: SettingsPtr,
  13. pub session: SessionWeakPtr,
  14. }
  15. impl Connector {
  16. /// Create a new connector with default network settings.
  17. pub fn new(settings: SettingsPtr, session: SessionWeakPtr) -> Self {
  18. Self { settings, session }
  19. }
  20. /// Establish an outbound connection.
  21. pub async fn connect(&self, connect_url: Url) -> Result<ChannelPtr> {
  22. let transport_name = TransportName::try_from(connect_url.clone())?;
  23. self.connect_channel(
  24. connect_url,
  25. transport_name,
  26. Duration::from_secs(self.settings.connect_timeout_seconds.into()),
  27. )
  28. .await
  29. }
  30. async fn connect_channel(
  31. &self,
  32. connect_url: Url,
  33. transport_name: TransportName,
  34. timeout: Duration,
  35. ) -> Result<Arc<Channel>> {
  36. macro_rules! connect {
  37. ($stream:expr, $transport:expr, $upgrade:expr) => {{
  38. if let Err(err) = $stream {
  39. error!("Setup for {} failed: {}", connect_url, err);
  40. return Err(Error::ConnectFailed)
  41. }
  42. let stream = $stream?.await;
  43. if let Err(err) = stream {
  44. error!("Connection to {} failed: {}", connect_url, err);
  45. return Err(Error::ConnectFailed)
  46. }
  47. let channel = match $upgrade {
  48. // session
  49. None => {
  50. Channel::new(Box::new(stream?), connect_url.clone(), self.session.clone())
  51. .await
  52. }
  53. Some(u) if u == "tls" => {
  54. let stream = $transport.upgrade_dialer(stream?)?.await;
  55. Channel::new(Box::new(stream?), connect_url, self.session.clone()).await
  56. }
  57. Some(u) => return Err(Error::UnsupportedTransportUpgrade(u)),
  58. };
  59. Ok(channel)
  60. }};
  61. }
  62. match transport_name {
  63. TransportName::Tcp(upgrade) => {
  64. let transport = TcpTransport::new(None, 1024);
  65. let stream = transport.dial(connect_url.clone(), Some(timeout));
  66. connect!(stream, transport, upgrade)
  67. }
  68. TransportName::Tor(upgrade) => {
  69. let socks5_url = Url::parse(
  70. &env::var("DARKFI_TOR_SOCKS5_URL")
  71. .unwrap_or_else(|_| "socks5://127.0.0.1:9050".to_string()),
  72. )?;
  73. let transport = TorTransport::new(socks5_url, None)?;
  74. let stream = transport.clone().dial(connect_url.clone(), None);
  75. connect!(stream, transport, upgrade)
  76. }
  77. _ => unimplemented!(),
  78. }
  79. }
  80. }