cashierd.rs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497
  1. use async_executor::Executor;
  2. use async_std::sync::{Arc, Mutex};
  3. use async_trait::async_trait;
  4. use clap::clap_app;
  5. use ff::Field;
  6. use log::{debug, warn};
  7. use rand::rngs::OsRng;
  8. use serde_json::{json, Value};
  9. use std::collections::HashMap;
  10. use std::path::PathBuf;
  11. use drk::{
  12. blockchain::Rocks,
  13. cli::{CashierdConfig, Config},
  14. client::Client,
  15. rpc::{
  16. jsonrpc::{error as jsonerr, response as jsonresp},
  17. jsonrpc::{ErrorCode::*, JsonRequest, JsonResult},
  18. rpcserver::{listen_and_serve, RequestHandler, RpcServerConfig},
  19. },
  20. serial::{deserialize, serialize},
  21. service::{bridge, bridge::Bridge},
  22. util::{expand_path, generate_id, join_config_path},
  23. wallet::{CashierDb, WalletDb},
  24. Error, Result,
  25. };
  26. fn handle_bridge_error(error_code: u32) -> Result<()> {
  27. match error_code {
  28. 1 => Err(Error::BridgeError("Not Supported Client".into())),
  29. _ => Err(Error::BridgeError("Unknown error_code".into())),
  30. }
  31. }
  32. #[derive(Clone)]
  33. struct Cashierd {
  34. config: CashierdConfig,
  35. bridge: Arc<Bridge>,
  36. cashier_wallet: Arc<CashierDb>,
  37. features: HashMap<String, String>,
  38. client: Arc<Mutex<Client>>,
  39. }
  40. #[async_trait]
  41. impl RequestHandler for Cashierd {
  42. // TODO: ServerError codes should be part of the lib.
  43. async fn handle_request(&self, req: JsonRequest) -> JsonResult {
  44. if req.params.as_array().is_none() {
  45. return JsonResult::Err(jsonerr(InvalidParams, None, req.id));
  46. }
  47. debug!(target: "RPC", "--> {}", serde_json::to_string(&req).unwrap());
  48. match req.method.as_str() {
  49. Some("deposit") => return self.deposit(req.id, req.params).await,
  50. Some("withdraw") => return self.withdraw(req.id, req.params).await,
  51. Some("features") => return self.features(req.id, req.params).await,
  52. Some(_) => {}
  53. None => {}
  54. };
  55. return JsonResult::Err(jsonerr(MethodNotFound, None, req.id));
  56. }
  57. }
  58. impl Cashierd {
  59. fn new(config_path: PathBuf) -> Result<Self> {
  60. let config: CashierdConfig = Config::<CashierdConfig>::load(config_path)?;
  61. let cashier_wallet = CashierDb::new(
  62. expand_path(&config.cashier_wallet_path.clone())?.as_path(),
  63. config.cashier_wallet_password.clone(),
  64. )?;
  65. let client_wallet = WalletDb::new(
  66. expand_path(&config.client_wallet_path.clone())?.as_path(),
  67. config.client_wallet_password.clone(),
  68. )?;
  69. let rocks = Rocks::new(expand_path(&config.database_path.clone())?.as_path())?;
  70. let client = Client::new(
  71. rocks,
  72. (
  73. config.gateway_protocol_url.parse()?,
  74. config.gateway_publisher_url.parse()?,
  75. ),
  76. (
  77. expand_path(&config.mint_params_path.clone())?,
  78. expand_path(&config.spend_params_path.clone())?,
  79. ),
  80. client_wallet.clone(),
  81. )?;
  82. let client = Arc::new(Mutex::new(client));
  83. let mut features = HashMap::new();
  84. for network in config.clone().networks {
  85. features.insert(network.name, network.blockchain);
  86. }
  87. let bridge = bridge::Bridge::new();
  88. Ok(Self {
  89. config: config.clone(),
  90. bridge,
  91. cashier_wallet,
  92. features,
  93. client: client.clone(),
  94. })
  95. }
  96. async fn resume_watch_deposit_keys(
  97. bridge: Arc<Bridge>,
  98. cashier_wallet: Arc<CashierDb>,
  99. features: HashMap<String, String>,
  100. ) -> Result<()> {
  101. for (network, _) in features.iter() {
  102. let keypairs_to_watch = cashier_wallet.get_deposit_token_keys_by_network(&network)?;
  103. for keypair in keypairs_to_watch {
  104. let bridge = bridge.clone();
  105. let bridge_subscribtion = bridge.subscribe().await;
  106. bridge_subscribtion
  107. .sender
  108. .send(bridge::BridgeRequests {
  109. network: network.to_owned(),
  110. payload: bridge::BridgeRequestsPayload::Watch(Some((keypair.0, keypair.1))),
  111. })
  112. .await?;
  113. }
  114. }
  115. Ok(())
  116. }
  117. async fn listen_for_receiving_coins(
  118. bridge: Arc<Bridge>,
  119. cashier_wallet: Arc<CashierDb>,
  120. recv_coin: async_channel::Receiver<(jubjub::SubgroupPoint, u64)>,
  121. ) -> Result<()> {
  122. let bridge_subscribtion = bridge.subscribe().await;
  123. // received drk coin
  124. let (drk_pub_key, amount) = recv_coin.recv().await?;
  125. debug!(target: "CASHIER DAEMON", "Receive coin with following address and amount: {}, {}"
  126. , drk_pub_key, amount);
  127. // get public key, and asset_id of the token
  128. let token = cashier_wallet.get_withdraw_token_public_key_by_dkey_public(&drk_pub_key)?;
  129. // send a request to bridge to send equivalent amount of
  130. // received drk coin to token publickey
  131. if let Some((addr, network, _asset_id)) = token {
  132. bridge_subscribtion
  133. .sender
  134. .send(bridge::BridgeRequests {
  135. network: network.to_string(),
  136. payload: bridge::BridgeRequestsPayload::Send(addr.clone(), amount),
  137. })
  138. .await?;
  139. // receive a response
  140. let res = bridge_subscribtion.receiver.recv().await?;
  141. let error_code = res.error as u32;
  142. if error_code == 0 {
  143. match res.payload {
  144. bridge::BridgeResponsePayload::Send => {
  145. // TODO Send the received coins to the main address
  146. cashier_wallet.confirm_withdraw_key_record(&addr, &network)?;
  147. }
  148. _ => {}
  149. }
  150. } else {
  151. return handle_bridge_error(error_code);
  152. }
  153. }
  154. Ok(())
  155. }
  156. async fn deposit(&self, id: Value, params: Value) -> JsonResult {
  157. debug!(target: "CASHIER DAEMON", "RECEIVED DEPOSIT REQUEST");
  158. let args: &Vec<serde_json::Value> = params.as_array().unwrap();
  159. if args.len() != 3 {
  160. return JsonResult::Err(jsonerr(InvalidParams, None, id));
  161. }
  162. let network = &args[0].as_str().unwrap();
  163. let network = network.to_string();
  164. let token_id = &args[1];
  165. let drk_pub_key = &args[2].as_str().unwrap();
  166. if !self.features.contains_key(&network.clone()) {
  167. return JsonResult::Err(jsonerr(
  168. InvalidParams,
  169. Some(format!("Cashier doesn't support this network: {}", network)),
  170. id,
  171. ));
  172. }
  173. let result: Result<String> = async {
  174. Self::check_token_id(&network, token_id.as_str().unwrap())?;
  175. let asset_id = generate_id(token_id)?;
  176. let drk_pub_key = bs58::decode(&drk_pub_key).into_vec()?;
  177. let drk_pub_key: jubjub::SubgroupPoint = deserialize(&drk_pub_key)?;
  178. // check if the drk public key is already exist
  179. let check = self
  180. .cashier_wallet
  181. .get_deposit_token_keys_by_dkey_public(&drk_pub_key, &network)?;
  182. let bridge = self.bridge.clone();
  183. let bridge_subscribtion = bridge.subscribe().await;
  184. if check.is_empty() {
  185. bridge_subscribtion
  186. .sender
  187. .send(bridge::BridgeRequests {
  188. network: network.clone(),
  189. payload: bridge::BridgeRequestsPayload::Watch(None),
  190. })
  191. .await?;
  192. } else {
  193. let keypair = check[0].to_owned();
  194. bridge_subscribtion
  195. .sender
  196. .send(bridge::BridgeRequests {
  197. network: network.clone(),
  198. payload: bridge::BridgeRequestsPayload::Watch(Some((keypair.0, keypair.1))),
  199. })
  200. .await?;
  201. }
  202. let bridge_res = bridge_subscribtion.receiver.recv().await?;
  203. let error_code = bridge_res.error as u32;
  204. if error_code != 0 {
  205. return handle_bridge_error(error_code).map(|_| String::new());
  206. }
  207. match bridge_res.payload {
  208. bridge::BridgeResponsePayload::Watch(token_priv, token_pub) => {
  209. // add pairings to db
  210. self.cashier_wallet.put_deposit_keys(
  211. &drk_pub_key,
  212. &token_priv,
  213. &serialize(&token_pub),
  214. &network,
  215. &asset_id,
  216. )?;
  217. return Ok(token_pub);
  218. }
  219. bridge::BridgeResponsePayload::Address(token_pub) => {
  220. return Ok(token_pub);
  221. }
  222. _ => Err(Error::BridgeError(
  223. "Receive unknown value from Subscription".into(),
  224. )),
  225. }
  226. }
  227. .await;
  228. match result {
  229. Ok(res) => JsonResult::Resp(jsonresp(json!(res), json!(id))),
  230. Err(err) => JsonResult::Err(jsonerr(InternalError, Some(err.to_string()), json!(id))),
  231. }
  232. }
  233. async fn withdraw(&self, id: Value, params: Value) -> JsonResult {
  234. debug!(target: "CASHIER DAEMON", "RECEIVED DEPOSIT REQUEST");
  235. let args: &Vec<serde_json::Value> = params.as_array().unwrap();
  236. if args.len() != 4 {
  237. return JsonResult::Err(jsonerr(InvalidParams, None, id));
  238. }
  239. let network = &args[0].as_str().unwrap();
  240. let network = network.to_string();
  241. let token = &args[1];
  242. let address = &args[2].as_str().unwrap();
  243. let _amount = &args[3];
  244. if !self.features.contains_key(&network.clone()) {
  245. return JsonResult::Err(jsonerr(
  246. InvalidParams,
  247. Some(format!("Cashier doesn't support this network: {}", network)),
  248. id,
  249. ));
  250. }
  251. let result: Result<String> = async {
  252. Self::check_token_id(&network, token.as_str().unwrap())?;
  253. let asset_id = generate_id(&token)?;
  254. let address = serialize(&address.to_string());
  255. let cashier_public: jubjub::SubgroupPoint;
  256. if let Some(addr) = self
  257. .cashier_wallet
  258. .get_withdraw_keys_by_token_public_key(&address, &network)?
  259. {
  260. cashier_public = addr.public;
  261. } else {
  262. let cashier_secret = jubjub::Fr::random(&mut OsRng);
  263. cashier_public =
  264. zcash_primitives::constants::SPENDING_KEY_GENERATOR * cashier_secret;
  265. self.cashier_wallet.put_withdraw_keys(
  266. &address,
  267. &cashier_public,
  268. &cashier_secret,
  269. &network.to_string(),
  270. &asset_id,
  271. )?;
  272. }
  273. let cashier_public_str = bs58::encode(serialize(&cashier_public)).into_string();
  274. Ok(cashier_public_str)
  275. }
  276. .await;
  277. match result {
  278. Ok(res) => JsonResult::Resp(jsonresp(json!(res), json!(id))),
  279. Err(err) => JsonResult::Err(jsonerr(InternalError, Some(err.to_string()), json!(id))),
  280. }
  281. }
  282. async fn features(&self, id: Value, _params: Value) -> JsonResult {
  283. JsonResult::Resp(jsonresp(json!(self.features), id))
  284. }
  285. fn check_token_id(network: &str, token_id: &str) -> Result<()> {
  286. match network {
  287. #[cfg(feature = "sol")]
  288. "sol" | "solana" => {
  289. if token_id != "So11111111111111111111111111111111111111112" {
  290. // This is supposed to be a token mint account now
  291. use drk::service::sol::account_is_initialized_mint;
  292. use drk::service::sol::SolFailed::BadSolAddress;
  293. use solana_sdk::pubkey::Pubkey;
  294. use std::str::FromStr;
  295. let pubkey = match Pubkey::from_str(token_id) {
  296. Ok(v) => v,
  297. Err(e) => return Err(Error::from(BadSolAddress(e.to_string()))),
  298. };
  299. // FIXME: Use network name from variable
  300. if !account_is_initialized_mint("devnet".to_string(), &pubkey) {
  301. return Err(Error::CashierInvalidTokenId(
  302. "Given address is not a valid token mint".into(),
  303. ));
  304. }
  305. }
  306. }
  307. #[cfg(feature = "btc")]
  308. "btc" | "bitcoin" => {
  309. // Handle bitcoin address here if needed
  310. }
  311. _ => {}
  312. }
  313. Ok(())
  314. }
  315. async fn start(&self, executor: Arc<Executor<'static>>) -> Result<()> {
  316. self.cashier_wallet.init_db().await?;
  317. for (feature_name, chain) in self.features.iter() {
  318. let bridge2 = self.bridge.clone();
  319. match feature_name.as_str() {
  320. #[cfg(feature = "sol")]
  321. "sol" | "solana" => {
  322. debug!(target: "CASHIER DAEMON", "Add sol network");
  323. use drk::service::SolClient;
  324. use solana_sdk::signer::keypair::Keypair;
  325. let main_keypair: Keypair;
  326. let main_keypairs = self.cashier_wallet.get_main_keys(&"sol".into())?;
  327. if main_keypairs.is_empty() {
  328. main_keypair = Keypair::new();
  329. } else {
  330. main_keypair = deserialize(&main_keypairs[0].0)?;
  331. }
  332. let sol_client = SolClient::new(serialize(&main_keypair), &chain).await?;
  333. bridge2.add_clients("sol".into(), sol_client).await?;
  334. }
  335. #[cfg(feature = "btc")]
  336. "btc" | "bitcoin" => {
  337. debug!(target: "CASHIER DAEMON", "Add btc network");
  338. let btc_endpoint: (bitcoin::network::constants::Network, String) = (
  339. bitcoin::network::constants::Network::Bitcoin,
  340. String::from("ssl://blockstream.info:993"),
  341. );
  342. use drk::service::btc::BtcClient;
  343. let _btc_client = BtcClient::new(btc_endpoint)?;
  344. // NOTE bitcoin is not implemented yet
  345. //
  346. // TODO check if there is main_keypair inside
  347. // cashierdb before generating new one
  348. //
  349. //bridge2.add_clients("btc".into(), btc_client).await?;
  350. }
  351. _ => {
  352. warn!("No feature enabled for {} network", feature_name);
  353. }
  354. }
  355. }
  356. let resume_watch_deposit_keys_task = executor.spawn(Self::resume_watch_deposit_keys(
  357. self.bridge.clone(),
  358. self.cashier_wallet.clone(),
  359. self.features.clone(),
  360. ));
  361. self.client.lock().await.start().await?;
  362. let (notify, recv_coin) = async_channel::unbounded::<(jubjub::SubgroupPoint, u64)>();
  363. let cashier_client_subscriber_task =
  364. smol::spawn(Client::connect_to_subscriber_from_cashier(
  365. self.client.clone(),
  366. executor.clone(),
  367. self.cashier_wallet.clone(),
  368. notify.clone(),
  369. ));
  370. let cashier_wallet = self.cashier_wallet.clone();
  371. let bridge = self.bridge.clone();
  372. let listen_for_receiving_coins_task = smol::spawn(async move {
  373. loop {
  374. Self::listen_for_receiving_coins(
  375. bridge.clone(),
  376. cashier_wallet.clone(),
  377. recv_coin.clone(),
  378. )
  379. .await
  380. .expect(" listen for receiving coins");
  381. }
  382. });
  383. let cfg = RpcServerConfig {
  384. socket_addr: self.config.rpc_listen_address.clone(),
  385. use_tls: self.config.serve_tls,
  386. identity_path: expand_path(&self.config.clone().tls_identity_path)?,
  387. identity_pass: self.config.tls_identity_password.clone(),
  388. };
  389. listen_and_serve(cfg, self.clone()).await?;
  390. resume_watch_deposit_keys_task.cancel().await;
  391. listen_for_receiving_coins_task.cancel().await;
  392. cashier_client_subscriber_task.cancel().await;
  393. Ok(())
  394. }
  395. }
  396. #[async_std::main]
  397. async fn main() -> Result<()> {
  398. let args = clap_app!(cashierd =>
  399. (@arg CONFIG: -c --config +takes_value "Sets a custom config file")
  400. (@arg verbose: -v --verbose "Increase verbosity")
  401. )
  402. .get_matches();
  403. let config_path = if args.is_present("CONFIG") {
  404. PathBuf::from(args.value_of("CONFIG").unwrap())
  405. } else {
  406. join_config_path(&PathBuf::from("cashierd.toml"))?
  407. };
  408. let loglevel = if args.is_present("verbose") {
  409. log::Level::Debug
  410. } else {
  411. log::Level::Info
  412. };
  413. simple_logger::init_with_level(loglevel)?;
  414. let ex = Arc::new(Executor::new());
  415. let cashierd = Cashierd::new(config_path)?;
  416. cashierd.start(ex.clone()).await
  417. }