jsonserver.rs 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261
  1. use crate::rpc::adapter::RpcAdapter;
  2. use crate::cli::cli_config;
  3. use crate::{Error, Result};
  4. use async_executor::Executor;
  5. use async_native_tls::TlsAcceptor;
  6. use async_std::sync::Mutex;
  7. use http_types::{Request, Response, StatusCode};
  8. use log::*;
  9. use smol::Async;
  10. use std::net::TcpListener;
  11. use std::sync::Arc;
  12. /// Listens for incoming connections and serves them.
  13. pub async fn listen(
  14. executor: Arc<Executor<'_>>,
  15. rpc: Arc<RpcInterface>,
  16. listener: Async<TcpListener>,
  17. tls: Option<TlsAcceptor>,
  18. ) -> Result<()> {
  19. // Format the full host address.
  20. let host = match &tls {
  21. None => format!("http://{}", listener.get_ref().local_addr()?),
  22. Some(_) => format!("https://{}", listener.get_ref().local_addr()?),
  23. };
  24. println!("Listening on {}", host);
  25. loop {
  26. // Accept the next connection.
  27. debug!(target: "rpc", "waiting for stream accept [START]");
  28. let (stream, _) = listener.accept().await?;
  29. debug!(target: "rpc", "stream accepted [END]");
  30. // Spawn a background task serving this connection.
  31. let task = match &tls {
  32. None => {
  33. let stream = async_dup::Arc::new(stream);
  34. let rpc = rpc.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)
  39. })
  40. .await
  41. {
  42. println!("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. println!("Failed to establish secure TLS connection: {:#?}", err);
  60. continue;
  61. }
  62. }
  63. }
  64. };
  65. // Detach the task to let it run in the background.
  66. task.detach();
  67. }
  68. }
  69. pub async fn start(
  70. executor: Arc<Executor<'_>>,
  71. config: Arc<cli_config::Config>,
  72. adapter: RpcAdapter,
  73. ) -> Result<()> {
  74. let rpc = RpcInterface::new(adapter)?;
  75. let rpc_url: std::net::SocketAddr = config.rpc_url.parse()?;
  76. let http = listen(
  77. executor.clone(),
  78. rpc.clone(),
  79. Async::<TcpListener>::bind(rpc_url)?,
  80. None,
  81. );
  82. let http_task = executor.spawn(http);
  83. *rpc.started.lock().await = true;
  84. rpc.wait_for_quit().await?;
  85. http_task.cancel().await;
  86. Ok(())
  87. }
  88. // json RPC server goes here
  89. #[allow(dead_code)]
  90. pub struct RpcInterface {
  91. pub started: Mutex<bool>,
  92. stop_send: async_channel::Sender<()>,
  93. stop_recv: async_channel::Receiver<()>,
  94. adapter: RpcAdapter,
  95. }
  96. impl RpcInterface {
  97. pub fn new(adapter: RpcAdapter) -> Result<Arc<Self>> {
  98. let (stop_send, stop_recv) = async_channel::unbounded::<()>();
  99. Ok(Arc::new(Self {
  100. //p2p,
  101. started: Mutex::new(false),
  102. stop_send,
  103. stop_recv,
  104. adapter,
  105. }))
  106. }
  107. pub async fn serve(self: Arc<Self>, mut req: Request) -> http_types::Result<Response> {
  108. info!("RPC serving {}", req.url());
  109. let request = req.body_string().await?;
  110. let io = self.handle_input().await?;
  111. let response = io
  112. .handle_request_sync(&request)
  113. .ok_or(Error::BadOperationType)?;
  114. let mut res = Response::new(StatusCode::Ok);
  115. res.insert_header("Content-Type", "text/plain");
  116. res.set_body(response);
  117. Ok(res)
  118. }
  119. pub async fn handle_input(self: Arc<Self>) -> Result<jsonrpc_core::IoHandler> {
  120. debug!(target: "rpc", "JsonRpcInterface::handle_input() [START]");
  121. let mut io = jsonrpc_core::IoHandler::new();
  122. io.add_sync_method("say_hello", |_| {
  123. Ok(jsonrpc_core::Value::String("Hello World!".into()))
  124. });
  125. let self1 = self.clone();
  126. io.add_method("get_key", move |_| {
  127. let self2 = self1.clone();
  128. async move {
  129. self2.adapter.get_key().await.expect("Failed to get key");
  130. Ok(jsonrpc_core::Value::String("Getting cashier key...".into()))
  131. }
  132. });
  133. let self1 = self.clone();
  134. io.add_method("get_cash_key", move |_| {
  135. let self2 = self1.clone();
  136. async move {
  137. self2
  138. .adapter
  139. .get_cash_key()
  140. .await
  141. .expect("Failed to get key");
  142. Ok(jsonrpc_core::Value::String("Getting cashier key...".into()))
  143. }
  144. });
  145. let self1 = self.clone();
  146. io.add_method("get_info", move |_| {
  147. let self2 = self1.clone();
  148. async move {
  149. self2.adapter.get_info().await;
  150. Ok(jsonrpc_core::Value::Null)
  151. }
  152. });
  153. let self1 = self.clone();
  154. io.add_method("stop", move |_| {
  155. let self2 = self1.clone();
  156. async move {
  157. self2.adapter.stop().await;
  158. Ok(jsonrpc_core::Value::Null)
  159. }
  160. });
  161. let self1 = self.clone();
  162. io.add_method("create_wallet", move |_| {
  163. let self2 = self1.clone();
  164. async move {
  165. println!(
  166. "Attempting wallet generation at path {:?}",
  167. self2.adapter.wallet.path
  168. );
  169. self2
  170. .adapter
  171. .init_db()
  172. .await
  173. .expect("Wallet generation failed");
  174. Ok(jsonrpc_core::Value::String("Created wallet".into()))
  175. }
  176. });
  177. let self1 = self.clone();
  178. io.add_method("key_gen", move |_| {
  179. let self2 = self1.clone();
  180. async move {
  181. println!("Key generation method called...");
  182. self2
  183. .adapter
  184. .key_gen()
  185. .await
  186. .expect("Failed to generate key");
  187. Ok(jsonrpc_core::Value::String(
  188. "Attempted key generation".into(),
  189. ))
  190. }
  191. });
  192. let self1 = self.clone();
  193. io.add_method("cash_key_gen", move |_| {
  194. let self2 = self1.clone();
  195. async move {
  196. println!("Key generation method called...");
  197. self2
  198. .adapter
  199. .cash_key_gen()
  200. .await
  201. .expect("Failed to generate key");
  202. Ok(jsonrpc_core::Value::String(
  203. "Attempted key generation".into(),
  204. ))
  205. }
  206. });
  207. let self1 = self.clone();
  208. io.add_method("test_wallet", move |_| {
  209. let self2 = self1.clone();
  210. async move {
  211. println!("Test wallet method called...");
  212. // use map err to convert from own error to jsonrpc
  213. // convert our error to string
  214. // use json to process error string
  215. self2.adapter.test_wallet().await.expect("Wallet test failed");
  216. Ok(jsonrpc_core::Value::String("Test wallet".into()))
  217. }
  218. });
  219. let self1 = self.clone();
  220. io.add_method("create_cashier_wallet", move |_| {
  221. let self2 = self1.clone();
  222. async move {
  223. println!("New wallet method called...");
  224. self2
  225. .adapter
  226. .init_cashier_db()
  227. .await.expect("Create wallet failed");
  228. println!("Wallet created at path {:?}", self2.adapter.wallet.path);
  229. Ok(jsonrpc_core::Value::String("Created cashier wallet".into()))
  230. }
  231. });
  232. debug!(target: "rpc", "JsonRpcInterface::handle_input() [END]");
  233. Ok(io)
  234. }
  235. pub async fn wait_for_quit(self: Arc<Self>) -> Result<()> {
  236. Ok(self.stop_recv.recv().await?)
  237. }
  238. }