rpcserver.rs 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182
  1. use async_std::{
  2. io::{ReadExt, WriteExt},
  3. sync::Arc,
  4. };
  5. use std::{env, fs};
  6. use async_trait::async_trait;
  7. use log::{debug, error, info};
  8. use url::Url;
  9. use super::jsonrpc::{JsonRequest, JsonResult};
  10. use crate::{
  11. net::transport::{
  12. TcpTransport, TorTransport, Transport, TransportListener, TransportName, TransportStream,
  13. },
  14. Error, Result,
  15. };
  16. #[async_trait]
  17. pub trait RequestHandler: Sync + Send {
  18. async fn handle_request(&self, req: JsonRequest) -> JsonResult;
  19. }
  20. async fn run_accept_loop(
  21. listener: Box<dyn TransportListener>,
  22. rh: Arc<impl RequestHandler + 'static>,
  23. ) -> Result<()> {
  24. // TODO can we spawn new task here ?
  25. while let Ok((stream, peer_addr)) = listener.next().await {
  26. info!(target: "JSON-RPC SERVER", "RPC Accepted connection {}", peer_addr);
  27. accept(stream, rh.clone()).await?;
  28. }
  29. Ok(())
  30. }
  31. async fn accept(
  32. mut stream: Box<dyn TransportStream>,
  33. rh: Arc<impl RequestHandler + 'static>,
  34. ) -> Result<()> {
  35. let mut buf = vec![0; 8192];
  36. loop {
  37. let n = match stream.read(&mut buf).await {
  38. Ok(n) if n == 0 => {
  39. info!(target: "JSON-RPC SERVER", "Closed connection");
  40. break
  41. }
  42. Ok(n) => n,
  43. Err(e) => {
  44. error!(target: "JSON-RPC SERVER", "Failed reading from socket: {}", e);
  45. info!(target: "JSON-RPC SERVER", "Closed connection");
  46. break
  47. }
  48. };
  49. let r: JsonRequest = match serde_json::from_slice(&buf[0..n]) {
  50. Ok(r) => {
  51. debug!(target: "JSON-RPC SERVER", "--> {}", String::from_utf8_lossy(&buf));
  52. r
  53. }
  54. Err(e) => {
  55. error!(target: "JSON-RPC SERVER", "Received invalid JSON: {:?}", e);
  56. info!(target: "JSON-RPC SERVER", "Closed connection");
  57. break
  58. }
  59. };
  60. let reply = rh.handle_request(r).await;
  61. let j = serde_json::to_string(&reply)?;
  62. debug!(target: "JSON-RPC SERVER", "<-- {}", j);
  63. if let Err(e) = stream.write_all(j.as_bytes()).await {
  64. error!(target: "JSON-RPC SERVER", "Failed writing to socket: {}", e);
  65. info!(target: "JSON-RPC SERVER", "Closed connection");
  66. break
  67. }
  68. }
  69. Ok(())
  70. }
  71. pub async fn listen_and_serve(
  72. accept_url: Url,
  73. rh: Arc<impl RequestHandler + 'static>,
  74. ) -> Result<()> {
  75. debug!(target: "JSON-RPC SERVER", "Trying to start listener on {}", accept_url);
  76. let transport_name = TransportName::try_from(accept_url.clone())?;
  77. match transport_name {
  78. TransportName::Tcp(upgrade) => {
  79. let transport = TcpTransport::new(None, 1024);
  80. let listener = transport.listen_on(accept_url.clone());
  81. if let Err(err) = listener {
  82. error!("TCP Setup failed: {}", err);
  83. return Err(Error::BindFailed(accept_url.clone().to_string()))
  84. }
  85. let listener = listener?.await;
  86. if let Err(err) = listener {
  87. error!("TCP Bind listener failed: {}", err);
  88. return Err(Error::BindFailed(accept_url.to_string()))
  89. }
  90. let listener = listener?;
  91. match upgrade {
  92. None => {
  93. run_accept_loop(Box::new(listener), rh).await?;
  94. }
  95. Some(u) if u == "tls" => {
  96. let tls_listener = transport.upgrade_listener(listener)?.await?;
  97. run_accept_loop(Box::new(tls_listener), rh).await?;
  98. }
  99. Some(u) => return Err(Error::UnsupportedTransportUpgrade(u)),
  100. }
  101. }
  102. TransportName::Tor(upgrade) => {
  103. let socks5_url = Url::parse(
  104. &env::var("DARKFI_TOR_SOCKS5_URL").unwrap_or("socks5://127.0.0.1:9050".to_string()),
  105. )?;
  106. let torc_url = Url::parse(
  107. &env::var("DARKFI_TOR_CONTROL_URL").unwrap_or("tcp://127.0.0.1:9051".to_string()),
  108. )?;
  109. let auth_cookie = env::var("DARKFI_TOR_COOKIE");
  110. if auth_cookie.is_err() {
  111. return Err(Error::TorError(
  112. "Please set the env var DARKFI_TOR_COOKIE to the configured tor cookie file. \
  113. For example: \
  114. \'export DARKFI_TOR_COOKIE=\"/var/lib/tor/control_auth_cookie\"\'"
  115. .to_string(),
  116. ))
  117. }
  118. let auth_cookie = auth_cookie.unwrap();
  119. let auth_cookie = hex::encode(&fs::read(auth_cookie).unwrap());
  120. let transport = TorTransport::new(socks5_url, Some((torc_url, auth_cookie)))?;
  121. // generate EHS pointing to local address
  122. let hurl = transport.create_ehs(accept_url.clone())?;
  123. info!("EHS TOR: {}", hurl.to_string());
  124. let listener = transport.clone().listen_on(accept_url.clone());
  125. if let Err(err) = listener {
  126. error!("TOR Setup failed: {}", err);
  127. return Err(Error::BindFailed(accept_url.clone().to_string()))
  128. }
  129. let listener = listener?.await;
  130. if let Err(err) = listener {
  131. error!("TOR Bind listener failed: {}", err);
  132. return Err(Error::BindFailed(accept_url.to_string()))
  133. }
  134. let listener = listener?;
  135. match upgrade {
  136. None => {
  137. run_accept_loop(Box::new(listener), rh).await?;
  138. }
  139. Some(u) if u == "tls" => {
  140. let tls_listener = transport.upgrade_listener(listener)?.await?;
  141. run_accept_loop(Box::new(tls_listener), rh).await?;
  142. }
  143. Some(u) => return Err(Error::UnsupportedTransportUpgrade(u)),
  144. }
  145. }
  146. _ => unimplemented!(),
  147. }
  148. Ok(())
  149. }