transport.rs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::time::Duration;
  19. use async_trait::async_trait;
  20. use smol::io::{AsyncRead, AsyncWrite};
  21. use url::Url;
  22. use crate::{Error, Result};
  23. /// TLS Upgrade Mechanism
  24. pub(crate) mod tls;
  25. #[cfg(feature = "p2p-transport-tcp")]
  26. /// TCP Transport
  27. pub(crate) mod tcp;
  28. #[cfg(feature = "p2p-transport-tor")]
  29. /// Tor transport
  30. pub(crate) mod tor;
  31. #[cfg(feature = "p2p-transport-nym")]
  32. /// Nym transport
  33. pub(crate) mod nym;
  34. #[cfg(feature = "p2p-transport-unix")]
  35. /// Unix socket transport
  36. pub(crate) mod unix;
  37. /// Dialer variants
  38. #[derive(Debug, Clone)]
  39. pub enum DialerVariant {
  40. #[cfg(feature = "p2p-transport-tcp")]
  41. /// Plain TCP
  42. Tcp(tcp::TcpDialer),
  43. #[cfg(feature = "p2p-transport-tcp")]
  44. /// TCP with TLS
  45. TcpTls(tcp::TcpDialer),
  46. #[cfg(feature = "p2p-transport-tor")]
  47. /// Tor
  48. Tor(tor::TorDialer),
  49. #[cfg(feature = "p2p-transport-tor")]
  50. /// Tor with TLS
  51. TorTls(tor::TorDialer),
  52. #[cfg(feature = "p2p-transport-nym")]
  53. /// Nym
  54. Nym(nym::NymDialer),
  55. #[cfg(feature = "p2p-transport-nym")]
  56. /// Nym with TLS
  57. NymTls(nym::NymDialer),
  58. #[cfg(feature = "p2p-transport-unix")]
  59. /// Unix socket
  60. Unix(unix::UnixDialer),
  61. }
  62. /// Listener variants
  63. #[derive(Debug, Clone)]
  64. pub enum ListenerVariant {
  65. #[cfg(feature = "p2p-transport-tcp")]
  66. /// Plain TCP
  67. Tcp(tcp::TcpListener),
  68. #[cfg(feature = "p2p-transport-tcp")]
  69. /// TCP with TLS
  70. TcpTls(tcp::TcpListener),
  71. #[cfg(feature = "p2p-transport-unix")]
  72. /// Unix socket
  73. Unix(unix::UnixListener),
  74. }
  75. /// A dialer that is able to transparently operate over arbitrary transports.
  76. pub struct Dialer {
  77. /// The endpoint to connect to
  78. endpoint: Url,
  79. /// The dialer variant (transport protocol)
  80. variant: DialerVariant,
  81. }
  82. macro_rules! enforce_hostport {
  83. ($endpoint:ident) => {
  84. if $endpoint.host_str().is_none() || $endpoint.port().is_none() {
  85. return Err(Error::InvalidDialerScheme)
  86. }
  87. };
  88. }
  89. macro_rules! enforce_abspath {
  90. ($endpoint:ident) => {
  91. if $endpoint.host_str().is_some() || $endpoint.port().is_some() {
  92. return Err(Error::InvalidDialerScheme)
  93. }
  94. if $endpoint.to_file_path().is_err() {
  95. return Err(Error::InvalidDialerScheme)
  96. }
  97. };
  98. }
  99. impl Dialer {
  100. /// Instantiate a new [`Dialer`] with the given [`Url`].
  101. pub async fn new(endpoint: Url) -> Result<Self> {
  102. match endpoint.scheme().to_lowercase().as_str() {
  103. #[cfg(feature = "p2p-transport-tcp")]
  104. "tcp" => {
  105. // Build a TCP dialer
  106. enforce_hostport!(endpoint);
  107. let variant = tcp::TcpDialer::new(None).await?;
  108. let variant = DialerVariant::Tcp(variant);
  109. Ok(Self { endpoint, variant })
  110. }
  111. #[cfg(feature = "p2p-transport-tcp")]
  112. "tcp+tls" => {
  113. // Build a TCP dialer wrapped with TLS
  114. enforce_hostport!(endpoint);
  115. let variant = tcp::TcpDialer::new(None).await?;
  116. let variant = DialerVariant::TcpTls(variant);
  117. Ok(Self { endpoint, variant })
  118. }
  119. #[cfg(feature = "p2p-transport-tor")]
  120. "tor" => {
  121. // Build a Tor dialer
  122. enforce_hostport!(endpoint);
  123. let variant = tor::TorDialer::new().await?;
  124. let variant = DialerVariant::Tor(variant);
  125. Ok(Self { endpoint, variant })
  126. }
  127. #[cfg(feature = "p2p-transport-tor")]
  128. "tor+tls" => {
  129. // Build a Tor dialer wrapped with TLS
  130. enforce_hostport!(endpoint);
  131. let variant = tor::TorDialer::new().await?;
  132. let variant = DialerVariant::TorTls(variant);
  133. Ok(Self { endpoint, variant })
  134. }
  135. #[cfg(feature = "p2p-transport-nym")]
  136. "nym" => {
  137. // Build a Nym dialer
  138. enforce_hostport!(endpoint);
  139. let variant = nym::NymDialer::new().await?;
  140. let variant = DialerVariant::Nym(variant);
  141. Ok(Self { endpoint, variant })
  142. }
  143. #[cfg(feature = "p2p-transport-nym")]
  144. "nym+tls" => {
  145. // Build a Nym dialer wrapped with TLS
  146. enforce_hostport!(endpoint);
  147. let variant = nym::NymDialer::new().await?;
  148. let variant = DialerVariant::NymTls(variant);
  149. Ok(Self { endpoint, variant })
  150. }
  151. #[cfg(feature = "p2p-transport-unix")]
  152. "unix" => {
  153. enforce_abspath!(endpoint);
  154. // Build a Unix socket dialer
  155. let variant = unix::UnixDialer::new().await?;
  156. let variant = DialerVariant::Unix(variant);
  157. Ok(Self { endpoint, variant })
  158. }
  159. x => Err(Error::UnsupportedTransport(x.to_string())),
  160. }
  161. }
  162. /// Dial an instantiated [`Dialer`]. This creates a connection and returns a stream.
  163. pub async fn dial(&self, timeout: Option<Duration>) -> Result<Box<dyn PtStream>> {
  164. match &self.variant {
  165. #[cfg(feature = "p2p-transport-tcp")]
  166. DialerVariant::Tcp(dialer) => {
  167. // NOTE: sockaddr here is an array, can contain both ipv4 and ipv6
  168. let sockaddr = self.endpoint.socket_addrs(|| None)?;
  169. let stream = dialer.do_dial(sockaddr[0], timeout).await?;
  170. Ok(Box::new(stream))
  171. }
  172. #[cfg(feature = "p2p-transport-tcp")]
  173. DialerVariant::TcpTls(dialer) => {
  174. let sockaddr = self.endpoint.socket_addrs(|| None)?;
  175. let stream = dialer.do_dial(sockaddr[0], timeout).await?;
  176. let tlsupgrade = tls::TlsUpgrade::new();
  177. let stream = tlsupgrade.upgrade_dialer_tls(stream).await?;
  178. Ok(Box::new(stream))
  179. }
  180. #[cfg(feature = "p2p-transport-tor")]
  181. DialerVariant::Tor(dialer) => {
  182. let host = self.endpoint.host_str().unwrap();
  183. let port = self.endpoint.port().unwrap();
  184. let stream = dialer.do_dial(host, port, timeout).await?;
  185. Ok(Box::new(stream))
  186. }
  187. #[cfg(feature = "p2p-transport-tor")]
  188. DialerVariant::TorTls(dialer) => {
  189. let host = self.endpoint.host_str().unwrap();
  190. let port = self.endpoint.port().unwrap();
  191. let stream = dialer.do_dial(host, port, timeout).await?;
  192. let tlsupgrade = tls::TlsUpgrade::new();
  193. let stream = tlsupgrade.upgrade_dialer_tls(stream).await?;
  194. Ok(Box::new(stream))
  195. }
  196. #[cfg(feature = "p2p-transport-nym")]
  197. DialerVariant::Nym(_dialer) => {
  198. todo!();
  199. }
  200. #[cfg(feature = "p2p-transport-nym")]
  201. DialerVariant::NymTls(_dialer) => {
  202. todo!();
  203. }
  204. #[cfg(feature = "p2p-transport-unix")]
  205. DialerVariant::Unix(dialer) => {
  206. let path = self.endpoint.to_file_path()?;
  207. let stream = dialer.do_dial(path).await?;
  208. Ok(Box::new(stream))
  209. }
  210. }
  211. }
  212. /// Return a reference to the `Dialer` endpoint
  213. pub fn endpoint(&self) -> &Url {
  214. &self.endpoint
  215. }
  216. }
  217. /// A listener that is able to transparently listen over arbitrary transports.
  218. pub struct Listener {
  219. /// The address to open the listener on
  220. endpoint: Url,
  221. /// The listener variant (transport protocol)
  222. variant: ListenerVariant,
  223. }
  224. impl Listener {
  225. /// Instantiate a new [`Listener`] with the given [`Url`].
  226. /// Must contain a scheme, host string, and a port.
  227. pub async fn new(endpoint: Url) -> Result<Self> {
  228. match endpoint.scheme().to_lowercase().as_str() {
  229. #[cfg(feature = "p2p-transport-tcp")]
  230. "tcp" => {
  231. // Build a TCP listener
  232. enforce_hostport!(endpoint);
  233. let variant = tcp::TcpListener::new().await?;
  234. let variant = ListenerVariant::Tcp(variant);
  235. Ok(Self { endpoint, variant })
  236. }
  237. #[cfg(feature = "p2p-transport-tcp")]
  238. "tcp+tls" => {
  239. // Build a TCP listener wrapped with TLS
  240. enforce_hostport!(endpoint);
  241. let variant = tcp::TcpListener::new().await?;
  242. let variant = ListenerVariant::TcpTls(variant);
  243. Ok(Self { endpoint, variant })
  244. }
  245. #[cfg(feature = "p2p-transport-unix")]
  246. "unix" => {
  247. enforce_abspath!(endpoint);
  248. let variant = unix::UnixListener::new().await?;
  249. let variant = ListenerVariant::Unix(variant);
  250. Ok(Self { endpoint, variant })
  251. }
  252. x => Err(Error::UnsupportedTransport(x.to_string())),
  253. }
  254. }
  255. /// Listen on an instantiated [`Listener`].
  256. /// This will open a socket and return the listener.
  257. pub async fn listen(&self) -> Result<Box<dyn PtListener>> {
  258. match &self.variant {
  259. #[cfg(feature = "p2p-transport-tcp")]
  260. ListenerVariant::Tcp(listener) => {
  261. let sockaddr = self.endpoint.socket_addrs(|| None)?;
  262. let l = listener.do_listen(sockaddr[0]).await?;
  263. Ok(Box::new(l))
  264. }
  265. #[cfg(feature = "p2p-transport-tcp")]
  266. ListenerVariant::TcpTls(listener) => {
  267. let sockaddr = self.endpoint.socket_addrs(|| None)?;
  268. let l = listener.do_listen(sockaddr[0]).await?;
  269. let tlsupgrade = tls::TlsUpgrade::new();
  270. let l = tlsupgrade.upgrade_listener_tcp_tls(l).await?;
  271. Ok(Box::new(l))
  272. }
  273. #[cfg(feature = "p2p-transport-unix")]
  274. ListenerVariant::Unix(listener) => {
  275. let path = self.endpoint.to_file_path()?;
  276. let l = listener.do_listen(&path.into()).await?;
  277. Ok(Box::new(l))
  278. }
  279. }
  280. }
  281. pub fn endpoint(&self) -> &Url {
  282. &self.endpoint
  283. }
  284. }
  285. /// Wrapper trait for async streams
  286. pub trait PtStream: AsyncRead + AsyncWrite + Unpin + Send {}
  287. #[cfg(feature = "p2p-transport-tcp")]
  288. impl PtStream for smol::net::TcpStream {}
  289. #[cfg(feature = "p2p-transport-tcp")]
  290. impl PtStream for async_rustls::TlsStream<smol::net::TcpStream> {}
  291. #[cfg(feature = "p2p-transport-tor")]
  292. impl PtStream for arti_client::DataStream {}
  293. #[cfg(feature = "p2p-transport-tor")]
  294. impl PtStream for async_rustls::TlsStream<arti_client::DataStream> {}
  295. #[cfg(feature = "p2p-transport-unix")]
  296. impl PtStream for smol::net::unix::UnixStream {}
  297. /// Wrapper trait for async listeners
  298. #[async_trait]
  299. pub trait PtListener: Send + Sync + Unpin {
  300. async fn next(&self) -> Result<(Box<dyn PtStream>, Url)>;
  301. }