jsonserver.rs 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. use crate::{Error, Result};
  2. use async_executor::Executor;
  3. use async_native_tls::TlsAcceptor;
  4. use async_std::sync::Mutex;
  5. use http_types::{Request, Response, StatusCode};
  6. use log::*;
  7. use smol::Async;
  8. use std::net::TcpListener;
  9. use std::sync::Arc;
  10. /// Listens for incoming connections and serves them.
  11. pub async fn listen(
  12. executor: Arc<Executor<'_>>,
  13. rpc: Arc<RpcInterface>,
  14. listener: Async<TcpListener>,
  15. tls: Option<TlsAcceptor>,
  16. io: Arc<jsonrpc_core::IoHandler>,
  17. ) -> Result<()> {
  18. // Format the full host address.
  19. let host = match &tls {
  20. None => format!("http://{}", listener.get_ref().local_addr()?),
  21. Some(_) => format!("https://{}", listener.get_ref().local_addr()?),
  22. };
  23. debug!(target: "RPC SERVER", "Listening on {}", host);
  24. loop {
  25. // Accept the next connection.
  26. debug!(target: "RPC SERVER", "waiting for stream accept [START]");
  27. let (stream, _) = listener.accept().await?;
  28. debug!(target: "RPC SERVER", "stream accepted [END]");
  29. // Spawn a background task serving this connection.
  30. let task = match &tls {
  31. None => {
  32. let stream = async_dup::Arc::new(stream);
  33. let rpc = rpc.clone();
  34. let io2 = io.clone();
  35. executor.spawn(async move {
  36. if let Err(err) = async_h1::accept(stream, move |req| {
  37. let rpc = rpc.clone();
  38. rpc.serve(req, io2.clone())
  39. })
  40. .await
  41. {
  42. debug!(target: "RPC SERVER", "Connection error: {:#?}", err);
  43. }
  44. })
  45. }
  46. Some(tls) => {
  47. // In case of HTTPS, establish a secure TLS connection first.
  48. match tls.accept(stream).await {
  49. Ok(stream) => {
  50. let _stream = async_dup::Arc::new(async_dup::Mutex::new(stream));
  51. executor.spawn(async move {
  52. /*if let Err(err) = async_h1::accept(stream, serve).await {
  53. println!("Connection error: {:#?}", err);
  54. }*/
  55. unimplemented!();
  56. })
  57. }
  58. Err(err) => {
  59. debug!(target: "RPC SERVER", "Failed to establish secure TLS connection: {:#?}", err);
  60. continue;
  61. }
  62. }
  63. }
  64. };
  65. task.await;
  66. }
  67. }
  68. pub async fn start(
  69. executor: Arc<Executor<'_>>,
  70. rpc_url: std::net::SocketAddr,
  71. io: Arc<jsonrpc_core::IoHandler>,
  72. ) -> Result<()> {
  73. let rpc = RpcInterface::new()?;
  74. let http = listen(
  75. executor.clone(),
  76. rpc.clone(),
  77. Async::<TcpListener>::bind(rpc_url)?,
  78. None,
  79. io,
  80. );
  81. executor.spawn(http).detach();
  82. *rpc.started.lock().await = true;
  83. Ok(())
  84. }
  85. #[allow(dead_code)]
  86. pub struct RpcInterface {
  87. pub started: Mutex<bool>,
  88. stop_send: async_channel::Sender<()>,
  89. stop_recv: async_channel::Receiver<()>,
  90. }
  91. impl RpcInterface {
  92. pub fn new() -> Result<Arc<Self>> {
  93. let (stop_send, stop_recv) = async_channel::unbounded::<()>();
  94. Ok(Arc::new(Self {
  95. started: Mutex::new(false),
  96. stop_send,
  97. stop_recv,
  98. }))
  99. }
  100. pub async fn serve(
  101. self: Arc<Self>,
  102. mut req: Request,
  103. io: Arc<jsonrpc_core::IoHandler>,
  104. ) -> http_types::Result<Response> {
  105. debug!(target: "RPC INTERFACE", "RPC serving {}", req.url());
  106. let request = req.body_string().await?;
  107. debug!(target: "RPC INTERFACE", "JsonRpcInterface::serve() [PROCESSING INPUT]");
  108. let response = io
  109. .handle_request_sync(&request)
  110. .ok_or(Error::BadOperationType)?;
  111. debug!(target: "RPC INTERFACE", "JsonRpcInterface::serve() [PROCESSED]");
  112. let mut res = Response::new(StatusCode::Ok);
  113. res.insert_header("Content-Type", "text/plain");
  114. res.set_body(response);
  115. Ok(res)
  116. }
  117. //pub async fn handle_input(self: Arc<Self>) -> Result<jsonrpc_core::IoHandler> {
  118. // debug!(target: "rpc", "JsonRpcInterface::handle_input() [START]");
  119. // let io = jsonrpc_core::IoHandler::new();
  120. // // this is where adapter is needed
  121. // //let io = self.adapter.clone().handle_input(io.clone())?;
  122. // debug!(target: "rpc", "JsonRpcInterface::handle_input() [END]");
  123. // Ok(io)
  124. //}
  125. pub async fn wait_for_quit(self: Arc<Self>) -> Result<()> {
  126. Ok(self.stop_recv.recv().await?)
  127. }
  128. }