jsonserver.rs 8.8 KB

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