server.rs 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158
  1. //! JSON-RPC server-side implementation.
  2. use async_std::sync::Arc;
  3. use async_trait::async_trait;
  4. use futures::{AsyncReadExt, AsyncWriteExt};
  5. use log::{debug, error, info, warn};
  6. use url::Url;
  7. use super::jsonrpc::{JsonRequest, JsonResult};
  8. use crate::{
  9. net::{
  10. transport::Transport, TcpTransport, TorTransport, TransportListener, TransportName,
  11. TransportStream, UnixTransport,
  12. },
  13. Error, Result,
  14. };
  15. /// Asynchronous trait implementing a handler for incoming JSON-RPC requests.
  16. /// Can be used by matching on methods and branching out to functions that
  17. /// handle respective methods.
  18. #[async_trait]
  19. pub trait RequestHandler: Sync + Send {
  20. async fn handle_request(&self, req: JsonRequest) -> JsonResult;
  21. }
  22. /// Internal accept function that runs inside a loop for accepting incoming
  23. /// JSON-RPC requests and passing them to the [`RequestHandler`].
  24. async fn accept(
  25. mut stream: Box<dyn TransportStream>,
  26. peer_addr: Url,
  27. rh: Arc<impl RequestHandler + 'static>,
  28. ) -> Result<()> {
  29. loop {
  30. // Nasty size
  31. let mut buf = vec![0; 8192 * 10];
  32. let n = match stream.read(&mut buf).await {
  33. Ok(n) if n == 0 => {
  34. debug!(target: "jsonrpc-server", "Closed connection for {}", peer_addr);
  35. break
  36. }
  37. Ok(n) => n,
  38. Err(e) => {
  39. error!("JSON-RPC server failed reading from {} socket: {}", peer_addr, e);
  40. debug!(target: "jsonrpc-server", "Closed connection for {}", peer_addr);
  41. break
  42. }
  43. };
  44. let r: JsonRequest = match serde_json::from_slice(&buf[0..n]) {
  45. Ok(r) => {
  46. debug!(target: "jsonrpc-server", "{} --> {}", peer_addr, String::from_utf8_lossy(&buf));
  47. r
  48. }
  49. Err(e) => {
  50. warn!("JSON-RPC server received invalid JSON from {}: {}", peer_addr, e);
  51. debug!(target: "jsonrpc-server", "Closed connection for {}", peer_addr);
  52. break
  53. }
  54. };
  55. let reply = rh.handle_request(r).await;
  56. let j = serde_json::to_string(&reply).unwrap();
  57. debug!(target: "jsonrpc-server", "{} <-- {}", peer_addr, j);
  58. if let Err(e) = stream.write_all(j.as_bytes()).await {
  59. error!("JSON-RPC server failed writing to {} socket: {}", peer_addr, e);
  60. debug!(target: "jsonrpc-server", "Closed connection for {}", peer_addr);
  61. break
  62. }
  63. }
  64. Ok(())
  65. }
  66. /// Wrapper function around [`accept()`] to take the incoming connection and
  67. /// pass it forward.
  68. async fn run_accept_loop(
  69. listener: Box<dyn TransportListener>,
  70. rh: Arc<impl RequestHandler + 'static>,
  71. ) -> Result<()> {
  72. while let Ok((stream, peer_addr)) = listener.next().await {
  73. info!("JSON-RPC server accepted connection from {}", peer_addr);
  74. accept(stream, peer_addr, rh.clone()).await?;
  75. }
  76. Ok(())
  77. }
  78. /// Start a JSON-RPC server bound to the given accept URL and use the given
  79. /// [`RequestHandler`] to handle incoming requests.
  80. pub async fn listen_and_serve(
  81. accept_url: Url,
  82. rh: Arc<impl RequestHandler + 'static>,
  83. ) -> Result<()> {
  84. debug!(target: "jsonrpc-server", "Trying to bind listener on {}", accept_url);
  85. macro_rules! accept {
  86. ($listener:expr, $transport:expr, $upgrade:expr) => {{
  87. if let Err(err) = $listener {
  88. error!("JSON-RPC server setup for {} failed: {}", accept_url, err);
  89. return Err(Error::BindFailed(accept_url.as_str().into()))
  90. }
  91. let listener = $listener?.await;
  92. if let Err(err) = listener {
  93. error!("JSON-RPC listener bind to {} failed: {}", accept_url, err);
  94. return Err(Error::BindFailed(accept_url.as_str().into()))
  95. }
  96. let listener = listener?;
  97. match $upgrade {
  98. None => {
  99. info!("JSON-RPC listener bound to {}", accept_url);
  100. run_accept_loop(Box::new(listener), rh).await?;
  101. }
  102. Some(u) if u == "tls" => {
  103. let tls_listener = $transport.upgrade_listener(listener)?.await?;
  104. info!("JSON-RPC listener bound to {}", accept_url);
  105. run_accept_loop(Box::new(tls_listener), rh).await?;
  106. }
  107. Some(u) => return Err(Error::UnsupportedTransportUpgrade(u)),
  108. }
  109. }};
  110. }
  111. let transport_name = TransportName::try_from(accept_url.clone())?;
  112. match transport_name {
  113. TransportName::Tcp(upgrade) => {
  114. let transport = TcpTransport::new(None, 1024);
  115. let listener = transport.listen_on(accept_url.clone());
  116. accept!(listener, transport, upgrade);
  117. }
  118. TransportName::Tor(upgrade) => {
  119. let (socks5_url, torc_url, auth_cookie) = TorTransport::get_listener_env()?;
  120. let auth_cookie = hex::encode(&std::fs::read(auth_cookie).unwrap());
  121. let transport = TorTransport::new(socks5_url, Some((torc_url, auth_cookie)))?;
  122. // Generate EHS pointing to local address
  123. let hurl = transport.create_ehs(accept_url.clone())?;
  124. info!("Created ephemeral hidden service: {}", hurl.to_string());
  125. let listener = transport.clone().listen_on(accept_url.clone());
  126. accept!(listener, transport, upgrade);
  127. }
  128. TransportName::Unix => {
  129. let transport = UnixTransport::new();
  130. let listener = transport.listen(accept_url.clone()).await;
  131. if let Err(err) = listener {
  132. error!("JSON-RPC Unix socket bind to {} failed: {}", accept_url, err);
  133. return Err(Error::BindFailed(accept_url.as_str().into()))
  134. }
  135. run_accept_loop(Box::new(listener?), rh).await?;
  136. }
  137. _ => unimplemented!(),
  138. }
  139. Ok(())
  140. }