server.rs 7.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204
  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. //! JSON-RPC server-side implementation.
  19. use async_std::sync::Arc;
  20. use async_trait::async_trait;
  21. use futures::{AsyncReadExt, AsyncWriteExt};
  22. use log::{debug, error, info, warn};
  23. use url::Url;
  24. use super::jsonrpc::{JsonRequest, JsonResult};
  25. use crate::{
  26. net::transport::{
  27. TcpTransport, TorTransport, Transport, TransportListener, TransportName, TransportStream,
  28. UnixTransport,
  29. },
  30. Error, Result,
  31. };
  32. /// Asynchronous trait implementing a handler for incoming JSON-RPC requests.
  33. /// Can be used by matching on methods and branching out to functions that
  34. /// handle respective methods.
  35. #[async_trait]
  36. pub trait RequestHandler: Sync + Send {
  37. async fn handle_request(&self, req: JsonRequest) -> JsonResult;
  38. }
  39. /// Internal accept function that runs inside a loop for accepting incoming
  40. /// JSON-RPC requests and passing them to the [`RequestHandler`].
  41. async fn accept(
  42. mut stream: Box<dyn TransportStream>,
  43. peer_addr: Url,
  44. rh: Arc<impl RequestHandler + 'static>,
  45. ) -> Result<()> {
  46. loop {
  47. // FIXME: Nasty size. 8M
  48. let mut buf = vec![0; 1024 * 8192];
  49. let n = match stream.read(&mut buf).await {
  50. Ok(n) if n == 0 => {
  51. debug!(target: "rpc::server", "Closed connection for {}", peer_addr);
  52. break
  53. }
  54. Ok(n) => n,
  55. Err(e) => {
  56. error!(target: "rpc::server", "JSON-RPC server failed reading from {} socket: {}", peer_addr, e);
  57. debug!(target: "rpc::server", "Closed connection for {}", peer_addr);
  58. break
  59. }
  60. };
  61. let r: JsonRequest = match serde_json::from_slice(&buf[0..n]) {
  62. Ok(r) => {
  63. debug!(target: "rpc::server", "{} --> {}", peer_addr, String::from_utf8_lossy(&buf));
  64. r
  65. }
  66. Err(e) => {
  67. warn!(target: "rpc::server", "JSON-RPC server received invalid JSON from {}: {}", peer_addr, e);
  68. debug!(target: "rpc::server", "Closed connection for {}", peer_addr);
  69. break
  70. }
  71. };
  72. let reply = rh.handle_request(r).await;
  73. match reply {
  74. JsonResult::Subscriber(sub) => {
  75. let subscription = sub.subscriber.subscribe().await;
  76. loop {
  77. // Listen subscription for notifications
  78. let notification = subscription.receive().await;
  79. // Push notification
  80. let j = serde_json::to_string(&notification).unwrap();
  81. debug!(target: "rpc::server", "{} <-- {}", peer_addr, j);
  82. if let Err(e) = stream.write_all(j.as_bytes()).await {
  83. error!(target: "rpc::server", "JSON-RPC server failed writing to {} socket: {}", peer_addr, e);
  84. debug!(target: "rpc::server", "Closed connection for {}", peer_addr);
  85. break
  86. }
  87. }
  88. subscription.unsubscribe().await;
  89. }
  90. _ => {
  91. let j = serde_json::to_string(&reply).unwrap();
  92. debug!(target: "rpc::server", "{} <-- {}", peer_addr, j);
  93. if let Err(e) = stream.write_all(j.as_bytes()).await {
  94. error!(target: "rpc::server", "JSON-RPC server failed writing to {} socket: {}", peer_addr, e);
  95. debug!(target: "rpc::server", "Closed connection for {}", peer_addr);
  96. break
  97. }
  98. }
  99. }
  100. }
  101. Ok(())
  102. }
  103. /// Wrapper function around [`accept()`] to take the incoming connection and
  104. /// pass it forward.
  105. async fn run_accept_loop(
  106. listener: Box<dyn TransportListener>,
  107. rh: Arc<impl RequestHandler + 'static>,
  108. ex: Arc<smol::Executor<'_>>,
  109. ) -> Result<()> {
  110. while let Ok((stream, peer_addr)) = listener.next().await {
  111. info!(target: "rpc::server", "JSON-RPC server accepted connection from {}", peer_addr);
  112. // Detaching requests handling
  113. let _rh = rh.clone();
  114. ex.spawn(async move {
  115. if let Err(e) = accept(stream, peer_addr.clone(), _rh).await {
  116. error!(target: "rpc::server", "JSON-RPC server error on handling request of {}: {}", peer_addr, e);
  117. }
  118. }).detach();
  119. }
  120. Ok(())
  121. }
  122. /// Start a JSON-RPC server bound to the given accept URL and use the given
  123. /// [`RequestHandler`] to handle incoming requests.
  124. pub async fn listen_and_serve(
  125. accept_url: Url,
  126. rh: Arc<impl RequestHandler + 'static>,
  127. ex: Arc<smol::Executor<'_>>,
  128. ) -> Result<()> {
  129. debug!(target: "rpc::server", "Trying to bind listener on {}", accept_url);
  130. macro_rules! accept {
  131. ($listener:expr, $transport:expr, $upgrade:expr) => {{
  132. if let Err(err) = $listener {
  133. error!(target: "rpc::server", "JSON-RPC server setup for {} failed: {}", accept_url, err);
  134. return Err(Error::BindFailed(accept_url.as_str().into()))
  135. }
  136. let listener = $listener?.await;
  137. if let Err(err) = listener {
  138. error!(target: "rpc::server", "JSON-RPC listener bind to {} failed: {}", accept_url, err);
  139. return Err(Error::BindFailed(accept_url.as_str().into()))
  140. }
  141. let listener = listener?;
  142. match $upgrade {
  143. None => {
  144. info!(target: "rpc::server", "JSON-RPC listener bound to {}", accept_url);
  145. run_accept_loop(Box::new(listener), rh, ex.clone()).await?;
  146. }
  147. Some(u) if u == "tls" => {
  148. let tls_listener = $transport.upgrade_listener(listener)?.await?;
  149. info!(target: "rpc::server", "JSON-RPC listener bound to {}", accept_url);
  150. run_accept_loop(Box::new(tls_listener), rh, ex.clone()).await?;
  151. }
  152. Some(u) => return Err(Error::UnsupportedTransportUpgrade(u)),
  153. }
  154. }};
  155. }
  156. let transport_name = TransportName::try_from(accept_url.clone())?;
  157. match transport_name {
  158. TransportName::Tcp(upgrade) => {
  159. let transport = TcpTransport::new(None, 1024);
  160. let listener = transport.listen_on(accept_url.clone());
  161. accept!(listener, transport, upgrade);
  162. }
  163. TransportName::Tor(upgrade) => {
  164. let (socks5_url, torc_url, auth_cookie) = TorTransport::get_listener_env()?;
  165. let auth_cookie = hex::encode(std::fs::read(auth_cookie).unwrap());
  166. let transport = TorTransport::new(socks5_url, Some((torc_url, auth_cookie)))?;
  167. // Generate EHS pointing to local address
  168. let hurl = transport.create_ehs(accept_url.clone())?;
  169. info!(target: "rpc::server", "Created ephemeral hidden service: {}", hurl.to_string());
  170. let listener = transport.clone().listen_on(accept_url.clone());
  171. accept!(listener, transport, upgrade);
  172. }
  173. TransportName::Unix => {
  174. let transport = UnixTransport::new();
  175. let listener = transport.listen_on(accept_url.clone());
  176. accept!(listener, transport, None);
  177. }
  178. _ => unimplemented!(),
  179. }
  180. Ok(())
  181. }