cashierd.rs 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497
  1. use drk::{
  2. blockchain::Rocks,
  3. cli::{CashierdConfig, Config},
  4. client::Client,
  5. rpc::{
  6. jsonrpc::{error as jsonerr, response as jsonresp},
  7. jsonrpc::{ErrorCode::*, JsonRequest, JsonResult},
  8. },
  9. serial::{deserialize, serialize},
  10. service::{bridge, bridge::Bridge},
  11. util::join_config_path,
  12. wallet::{CashierDb, WalletDb},
  13. Error, Result,
  14. };
  15. use clap::clap_app;
  16. use log::*;
  17. use serde::Serialize;
  18. use serde_json::{json, Value};
  19. use simplelog::{
  20. CombinedLogger, Config as SimLogConfig, ConfigBuilder, LevelFilter, TermLogger, TerminalMode,
  21. WriteLogger,
  22. };
  23. use tokio::io::{AsyncReadExt, AsyncWriteExt};
  24. use tokio::net::TcpListener;
  25. use async_executor::Executor;
  26. use easy_parallel::Parallel;
  27. use ff::Field;
  28. use rand::rngs::OsRng;
  29. use async_std::sync::{Arc, Mutex};
  30. use sha2::{Digest, Sha256};
  31. use std::path::PathBuf;
  32. #[derive(Debug, Clone, Serialize)]
  33. struct Features {
  34. networks: Vec<String>,
  35. }
  36. impl Features {
  37. fn new() -> Self {
  38. let mut networks = Vec::new();
  39. networks.push("solana".to_string());
  40. networks.push("bitcoin".to_string());
  41. Self { networks }
  42. }
  43. }
  44. #[derive(Clone)]
  45. struct Cashierd {
  46. verbose: bool,
  47. config: CashierdConfig,
  48. client_wallet: Arc<WalletDb>,
  49. cashier_wallet: Arc<CashierDb>,
  50. features: Features,
  51. client: Arc<Mutex<Client>>,
  52. }
  53. impl Cashierd {
  54. fn new(verbose: bool, config_path: PathBuf) -> Result<Self> {
  55. let mint_params_path = join_config_path(&PathBuf::from("cashier_mint.params"))?;
  56. let spend_params_path = join_config_path(&PathBuf::from("cashier_spend.params"))?;
  57. let config: CashierdConfig = Config::<CashierdConfig>::load(config_path)?;
  58. let cashier_wallet_path = join_config_path(&PathBuf::from("cashier_wallet.db"))?;
  59. let client_wallet_path = join_config_path(&PathBuf::from("cashier_client_wallet.db"))?;
  60. let cashier_wallet = CashierDb::new(&cashier_wallet_path, config.password.clone())?;
  61. let client_wallet = WalletDb::new(&client_wallet_path.clone(), config.password.clone())?;
  62. let database_path = join_config_path(&PathBuf::from("cashier_database.db"))?;
  63. let rocks = Rocks::new(&database_path)?;
  64. let client = Client::new(
  65. rocks,
  66. (
  67. config.gateway_url.parse()?,
  68. config.gateway_subscriber_url.parse()?,
  69. ),
  70. (mint_params_path, spend_params_path),
  71. client_wallet.clone(),
  72. )?;
  73. let client = Arc::new(Mutex::new(client));
  74. let features = Features::new();
  75. Ok(Self {
  76. verbose,
  77. config: config.clone(),
  78. cashier_wallet,
  79. client_wallet,
  80. features,
  81. client: client.clone(),
  82. })
  83. }
  84. async fn start(&self, executor: Arc<Executor<'_>>) -> Result<()> {
  85. self.cashier_wallet.init_db()?;
  86. let bridge = Bridge::new();
  87. self.client.lock().await.start().await?;
  88. let (notify, recv_coin) = async_channel::unbounded::<(jubjub::SubgroupPoint, u64)>();
  89. let cashier_client_subscriber_task =
  90. executor.spawn(Client::connect_to_subscriber_from_cashier(
  91. self.client.clone(),
  92. executor.clone(),
  93. self.cashier_wallet.clone(),
  94. notify.clone(),
  95. ));
  96. let cashier_wallet = self.cashier_wallet.clone();
  97. let ex = executor.clone();
  98. let listen_for_receiving_coins_task = executor.spawn(async move {
  99. loop {
  100. Self::listen_for_receiving_coins(
  101. ex.clone(),
  102. bridge.clone(),
  103. cashier_wallet.clone(),
  104. recv_coin.clone(),
  105. )
  106. .await
  107. .expect(" listen for receiving coins");
  108. }
  109. });
  110. let rpc_url = self.config.rpc_url.clone();
  111. run_rpc_server(self.clone(), rpc_url).await?;
  112. listen_for_receiving_coins_task.cancel().await;
  113. cashier_client_subscriber_task.cancel().await;
  114. Ok(())
  115. }
  116. async fn listen_for_receiving_coins(
  117. ex: Arc<Executor<'_>>,
  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(ex.clone()).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, asset_id)) = token {
  132. bridge_subscribtion
  133. .sender
  134. .send(bridge::BridgeRequests {
  135. asset_id,
  136. payload: bridge::BridgeRequestsPayload::SendRequest(addr.clone(), amount),
  137. })
  138. .await?;
  139. // receive a response
  140. let res = bridge_subscribtion.receiver.recv().await?;
  141. if res.error == 0 {
  142. match res.payload {
  143. bridge::BridgeResponsePayload::SendResponse => {
  144. // TODO Send the received coins to the main address
  145. cashier_wallet.confirm_withdraw_key_record(&addr, &asset_id)?;
  146. }
  147. _ => {}
  148. }
  149. }
  150. }
  151. Ok(())
  152. }
  153. async fn handle_request(self, req: JsonRequest) -> JsonResult {
  154. if req.params.as_array().is_none() {
  155. return JsonResult::Err(jsonerr(InvalidParams, None, req.id));
  156. }
  157. debug!(target: "RPC", "--> {:#?}", serde_json::to_string(&req).unwrap());
  158. match req.method.as_str() {
  159. Some("deposit") => return self.deposit(req.id, req.params).await,
  160. Some("withdraw") => return self.withdraw(req.id, req.params).await,
  161. Some("features") => return self.features(req.id, req.params).await,
  162. Some(_) => {}
  163. None => {}
  164. };
  165. return JsonResult::Err(jsonerr(MethodNotFound, None, req.id));
  166. }
  167. async fn deposit(self, id: Value, params: Value) -> JsonResult {
  168. debug!(target: "CASHIER", "RECEIVED DEPOSIT REQUEST");
  169. if params.as_array().is_none() {
  170. return JsonResult::Err(jsonerr(InvalidParams, None, id));
  171. }
  172. let args = params.as_array().unwrap();
  173. let _ntwk = &args[0];
  174. let tkn = &args[1];
  175. let pk = &args[2];
  176. debug!(target: "CASHIER", "PROCESSING INPUT");
  177. // TODO: proper error handling
  178. let token_id = Self::parse_id(tkn).unwrap();
  179. if pk.as_str().is_none() {
  180. return JsonResult::Err(jsonerr(InvalidParams, None, id));
  181. }
  182. let pk_str = pk.as_str().unwrap();
  183. let pk_58 = bs58::decode(pk_str).into_vec().unwrap();
  184. let pubkey: jubjub::SubgroupPoint = deserialize(&pk_58).unwrap();
  185. //// TODO: Sanity check.
  186. let _check = self
  187. .clone()
  188. .cashier_wallet
  189. .get_deposit_token_keys_by_dkey_public(&pubkey, &token_id);
  190. // TODO: implement bridge communication
  191. // this just returns the user public key
  192. let pubkey = bs58::encode(serialize(&pubkey)).into_string();
  193. debug!(target: "CASHIER", "ATTEMPING REPLY");
  194. JsonResult::Resp(jsonresp(json!(pubkey), json!(id)))
  195. }
  196. // here we hash the alphanumeric token ID. if it fails, we change the last 4 bytes and hash it
  197. // again, and keep repeating until it works.
  198. fn parse_id(token: &Value) -> Result<jubjub::Fr> {
  199. let tkn_str = token.as_str().unwrap();
  200. if bs58::decode(tkn_str).into_vec().is_err() {
  201. // TODO: make this an error
  202. debug!(target: "CASHIER", "COULD NOT DECODE STR");
  203. }
  204. let mut data = bs58::decode(tkn_str).into_vec().unwrap();
  205. let token_id = deserialize::<jubjub::Fr>(&data);
  206. if token_id.is_err() {
  207. let mut counter = 0;
  208. loop {
  209. data.truncate(28);
  210. let serialized_counter = serialize(&counter);
  211. data.extend(serialized_counter.iter());
  212. let mut hasher = Sha256::new();
  213. hasher.update(&data);
  214. let hash = hasher.finalize();
  215. let token_id = deserialize::<jubjub::Fr>(&hash);
  216. if token_id.is_err() {
  217. counter += 1;
  218. continue;
  219. }
  220. debug!(target: "CASHIER", "DESERIALIZATION SUCCESSFUL");
  221. let tkn = token_id.unwrap();
  222. return Ok(tkn);
  223. }
  224. }
  225. unreachable!();
  226. }
  227. async fn withdraw(self, id: Value, params: Value) -> JsonResult {
  228. debug!(target: "CASHIER DAEMON", "RECEIVED DEPOSIT REQUEST");
  229. // TODO Cashier checks if they support the network, and if so,
  230. // return adeposit address.
  231. let result: Result<String> = async {
  232. let args: &Vec<serde_json::Value>;
  233. if let Some(ar) = params.as_array() {
  234. args = ar;
  235. } else {
  236. return Err(Error::ParseFailed("Unable to parse rpc params to array"));
  237. }
  238. let _network = &args[0];
  239. let token = &args[1];
  240. let address = &args[2];
  241. let _amount = &args[3];
  242. let asset_id = Self::parse_id(&token)?;
  243. let address = serialize(&address.to_string());
  244. let cashier_public: jubjub::SubgroupPoint;
  245. if let Some(addr) = self
  246. .cashier_wallet
  247. .get_withdraw_keys_by_token_public_key(&address, &asset_id)?
  248. {
  249. cashier_public = addr.public;
  250. } else {
  251. let cashier_secret = jubjub::Fr::random(&mut OsRng);
  252. cashier_public =
  253. zcash_primitives::constants::SPENDING_KEY_GENERATOR * cashier_secret;
  254. self.cashier_wallet.put_withdraw_keys(
  255. &address,
  256. &cashier_public,
  257. &cashier_secret,
  258. &asset_id,
  259. )?;
  260. }
  261. let cashier_public_str = bs58::encode(serialize(&cashier_public)).into_string();
  262. Ok(cashier_public_str)
  263. }
  264. .await;
  265. match result {
  266. Ok(res) => JsonResult::Resp(jsonresp(json!(res), json!(id))),
  267. Err(err) => JsonResult::Err(jsonerr(InternalError, Some(err.to_string()), json!(id))),
  268. }
  269. }
  270. async fn features(self, id: Value, _params: Value) -> JsonResult {
  271. JsonResult::Resp(jsonresp(json!(self.features), id))
  272. }
  273. }
  274. async fn run_rpc_server(cashierd: Cashierd, rpc_url: String) -> Result<()> {
  275. let listener = TcpListener::bind(rpc_url.clone()).await?;
  276. debug!(target: "RPC SERVER", "Listening on {}", rpc_url);
  277. loop {
  278. debug!(target: "RPC SERVER", "waiting for client");
  279. let (mut socket, _) = listener.accept().await?;
  280. debug!(target: "RPC SERVER", "accepted client");
  281. let cashierd = cashierd.clone();
  282. tokio::spawn(async move {
  283. let mut buf = [0; 2048];
  284. loop {
  285. let n = match socket.read(&mut buf).await {
  286. Ok(n) if n == 0 => {
  287. debug!(target: "RPC SERVER", "closed connection");
  288. return;
  289. }
  290. Ok(n) => n,
  291. Err(e) => {
  292. debug!(target: "RPC SERVER", "failed to read from socket; err = {:?}", e);
  293. return;
  294. }
  295. };
  296. let r: JsonRequest = match serde_json::from_slice(&buf[0..n]) {
  297. Ok(r) => r,
  298. Err(e) => {
  299. debug!(target: "RPC SERVER", "received invalid json; err = {:?}", e);
  300. return;
  301. }
  302. };
  303. let reply = cashierd.clone().handle_request(r).await;
  304. let j = serde_json::to_string(&reply).unwrap();
  305. debug!(target: "RPC", "<-- {:#?}", j);
  306. // Write the data back
  307. if let Err(e) = socket.write_all(j.as_bytes()).await {
  308. debug!(target: "RPC SERVER", "failed to write to socket; err = {:?}", e);
  309. return;
  310. }
  311. }
  312. });
  313. }
  314. }
  315. #[tokio::main]
  316. async fn main() -> Result<()> {
  317. let args = clap_app!(cashierd =>
  318. (@arg CONFIG: -c --config +takes_value "Sets a custom config file")
  319. (@arg verbose: -v --verbose "Increase verbosity")
  320. )
  321. .get_matches();
  322. let config_path: PathBuf;
  323. if args.is_present("CONFIG") {
  324. config_path = PathBuf::from(args.value_of("CONFIG").unwrap());
  325. } else {
  326. config_path = join_config_path(&PathBuf::from("cashierd.toml"))?;
  327. }
  328. let cashierd = Cashierd::new(args.clone().is_present("verbose"), config_path)?;
  329. let logger_config = ConfigBuilder::new().set_time_format_str("%T%.6f").build();
  330. let debug_level = if args.is_present("verbose") {
  331. LevelFilter::Debug
  332. } else {
  333. LevelFilter::Off
  334. };
  335. let log_path = cashierd.clone().config.log_path;
  336. CombinedLogger::init(vec![
  337. TermLogger::new(debug_level, logger_config, TerminalMode::Mixed).unwrap(),
  338. WriteLogger::new(
  339. LevelFilter::Debug,
  340. SimLogConfig::default(),
  341. std::fs::File::create(log_path).unwrap(),
  342. ),
  343. ])
  344. .unwrap();
  345. let ex = Arc::new(Executor::new());
  346. let ex2 = ex.clone();
  347. let (signal, shutdown) = async_channel::unbounded::<()>();
  348. let cashierd2 = cashierd.clone();
  349. let (_, _result) = Parallel::new()
  350. // Run four executor threads.
  351. .each(0..3, |_| smol::future::block_on(ex.run(shutdown.recv())))
  352. // Run the main future on the current thread.
  353. .finish(|| {
  354. smol::future::block_on(async move {
  355. cashierd2.start(ex2).await?;
  356. drop(signal);
  357. Ok::<(), Error>(())
  358. })
  359. });
  360. Ok(())
  361. }
  362. #[cfg(test)]
  363. mod tests {
  364. use drk::serial::{deserialize, serialize};
  365. use sha2::{Digest, Sha256};
  366. #[test]
  367. fn test_jubjub_parsing() {
  368. // 1. counter = 0
  369. // 2. serialized_counter = serialize(counter)
  370. // 3. asset_id_data = hash(data + serialized_counter)
  371. // 4. asset_id = deserialize(asset_id_data)
  372. // 5. test parse
  373. // 6. loop
  374. let tkn_str = "EPjFWdd5AufqSSqeM2qN1xzybapC8G4wEGGkZwyTDt1v";
  375. println!("{}", tkn_str);
  376. if bs58::decode(tkn_str).into_vec().is_err() {
  377. println!("Could not decode str into vec");
  378. }
  379. let mut data = bs58::decode(tkn_str).into_vec().unwrap();
  380. println!("{:?}", data);
  381. let mut hasher = Sha256::new();
  382. hasher.update(&data);
  383. let hash = hasher.finalize();
  384. let token_id = deserialize::<jubjub::Fr>(&hash);
  385. println!("{:?}", token_id);
  386. let mut counter = 0;
  387. if token_id.is_err() {
  388. println!("could not deserialize tkn 58");
  389. loop {
  390. println!("TOKEN IS NONE. COMMENCING LOOP");
  391. counter += 1;
  392. println!("LOOP NUMBER {}", counter);
  393. println!("{:?}", data.len());
  394. data.truncate(28);
  395. let serialized_counter = serialize(&counter);
  396. println!("{:?}", serialized_counter);
  397. data.extend(serialized_counter.iter());
  398. println!("{:?}", data.len());
  399. let mut hasher = Sha256::new();
  400. hasher.update(&data);
  401. let hash = hasher.finalize();
  402. let token_id = deserialize::<jubjub::Fr>(&hash);
  403. println!("{:?}", token_id);
  404. if token_id.is_err() {
  405. continue;
  406. }
  407. if counter > 10 {
  408. break;
  409. }
  410. println!("deserialization successful");
  411. token_id.unwrap();
  412. break;
  413. }
  414. };
  415. }
  416. }