Просмотр исходного кода

cashierd: finish refactoring of start function

ghassmo 4 лет назад
Родитель
Сommit
1933575da6
1 измененных файлов с 96 добавлено и 49 удалено
  1. 96 49
      src/bin/cashierd2.rs

+ 96 - 49
src/bin/cashierd2.rs

@@ -1,35 +1,35 @@
-use async_std::sync::Arc;
-use log::*;
-use std::path::PathBuf;
+use drk::{
+    blockchain::Rocks,
+    cli::{CashierdConfig, Config},
+    client::Client,
+    rpc::{
+        jsonrpc::{error as jsonerr, response as jsonresp},
+        jsonrpc::{ErrorCode::*, JsonRequest, JsonResult},
+    },
+    serial::{deserialize, serialize},
+    service::bridge,
+    util::join_config_path,
+    wallet::{CashierDb, WalletDb},
+    Error, Result,
+};
 
 use clap::clap_app;
+use log::*;
 use serde::Serialize;
 use serde_json::{json, Value};
 use simplelog::{
     CombinedLogger, Config as SimLogConfig, ConfigBuilder, LevelFilter, TermLogger, TerminalMode,
     WriteLogger,
 };
-use std::net::SocketAddr;
 use tokio::io::{AsyncReadExt, AsyncWriteExt};
 use tokio::net::TcpListener;
 
 use async_executor::Executor;
 use easy_parallel::Parallel;
 
-use drk::{
-    cli::{CashierdConfig, Config},
-    rpc::{
-        jsonrpc::{error as jsonerr, response as jsonresp},
-        jsonrpc::{ErrorCode::*, JsonRequest, JsonResult},
-    },
-    serial::{deserialize, serialize},
-    service::{bridge, CashierService},
-    util::join_config_path,
-    wallet::{CashierDb, WalletDb},
-    Error, Result,
-};
-
+use async_std::sync::{Arc, Mutex};
 use ff::PrimeField;
+use std::path::PathBuf;
 
 #[derive(Debug, Clone, Serialize)]
 struct Features {
@@ -52,14 +52,16 @@ struct Cashierd {
     client_wallet: Arc<WalletDb>,
     cashier_wallet: Arc<CashierDb>,
     features: Features,
-    // clientdb:
-    // mint_params:
-    // spend_params:
+    client: Arc<Mutex<Client>>,
 }
 
 impl Cashierd {
     fn new(verbose: bool, config_path: PathBuf) -> Result<Self> {
+        let mint_params_path = join_config_path(&PathBuf::from("cashier_mint.params"))?;
+        let spend_params_path = join_config_path(&PathBuf::from("cashier_spend.params"))?;
+
         let config: CashierdConfig = Config::<CashierdConfig>::load(config_path)?;
+
         let cashier_wallet = CashierDb::new(
             &PathBuf::from(config.cashierdb_path.clone()),
             config.password.clone(),
@@ -68,45 +70,91 @@ impl Cashierd {
             &PathBuf::from(config.cashierdb_path.clone()),
             config.password.clone(),
         )?;
+
+        let rocks = Rocks::new(&PathBuf::from(&config.cashierdb_path))?;
+
+        let client = Client::new(
+            rocks,
+            (
+                config.gateway_url.parse()?,
+                config.gateway_subscriber_url.parse()?,
+            ),
+            (mint_params_path, spend_params_path),
+            client_wallet.clone(),
+        )?;
+
+        let client = Arc::new(Mutex::new(client));
+
         let features = Features::new();
 
         Ok(Self {
             verbose,
-            config,
+            config: config.clone(),
             cashier_wallet,
             client_wallet,
             features,
+            client: client.clone(),
         })
     }
 
-    async fn start(self, executor: Arc<Executor<'_>>, config: CashierdConfig) -> Result<()> {
-        let ex = executor.clone();
-        let accept_addr: SocketAddr = config.accept_url.parse()?;
+    async fn start(self, executor: Arc<Executor<'_>>) -> Result<()> {
+        //// TODO: pass vector of assets
 
-        let gateway_addr: SocketAddr = config.gateway_url.parse()?;
+        self.cashier_wallet.init_db()?;
 
-        let database_path = PathBuf::from(config.cashierdb_path);
+        let bridge = bridge::Bridge::new();
 
-        let mint_params_path = join_config_path(&PathBuf::from("cashier_mint.params"))?;
-        let spend_params_path = join_config_path(&PathBuf::from("cashier_spend.params"))?;
+        self.client.lock().await.start().await?;
 
-        let mut cashier = CashierService::new(
-            accept_addr,
-            self.cashier_wallet.clone(),
-            self.client_wallet.clone(),
-            database_path,
-            (gateway_addr, "127.0.0.1:4444".parse()?),
-            (mint_params_path, spend_params_path),
-        )
-        .await?;
+        let (notify, recv_coin) = async_channel::unbounded::<(jubjub::SubgroupPoint, u64)>();
+
+        let _cashier_client_subscriber_task =
+            executor.spawn(Client::connect_to_subscriber_from_cashier(
+                self.client.clone(),
+                executor.clone(),
+                self.cashier_wallet.clone(),
+                notify.clone(),
+            ));
+
+        let cashier_wallet = self.cashier_wallet.clone();
+
+        let ex = executor.clone();
+        executor.spawn(async move {
+            loop {
+                let bridge = bridge.clone();
+                let bridge_subscribtion  = bridge.subscribe(ex.clone()).await;
+                let (pub_key, amount) = recv_coin.recv().await.expect("Receive Own Coin");
+                debug!(target: "CASHIER DAEMON", "Receive coin with following address and amount: {}, {}", pub_key, amount);
+                let coin = cashier_wallet.get_withdraw_coin_public_key_by_dkey_public(&pub_key)
+                    .expect("Get coin_key by pub_key");
+                if let Some((addr, asset_id)) =  coin {
+                    // send equivalent amount of coin to this address
+                    bridge_subscribtion.sender.send(
+                        bridge::BridgeRequests {
+                            asset_id,
+                            payload: bridge::BridgeRequestsPayload::SendRequest(addr.clone(), amount)
+                        }
+                    ).await.expect("send request to bridge");
+
+                    let res = bridge_subscribtion.receiver.recv().await.expect("bridge resonse");
+
+                    if res.error == 0 {
+                        match res.payload {
+                            bridge::BridgeResponsePayload::SendResponse => {
+                                // TODO Send the received coins to the main address
+                                cashier_wallet.confirm_withdraw_key_record(&addr, &serialize(&1) )
+                                    .expect("Confirm withdraw key record");
+                            }
+                            _ => {}
+                        }
+
+                    }
 
-        //// TODO: make this a vector of accepted assets
-        //let asset = Asset::new("btc".to_string());
-        //// TODO: this should be done by the user
-        //let asset_id = deserialize(&asset.id)?;
 
-        //// TODO: pass vector of assets into cashier.start()
-        //cashier.start(ex.clone(), asset_id).await?;
+                }
+
+            }
+        }).await;
 
         Ok(())
     }
@@ -183,10 +231,10 @@ impl Cashierd {
 
         let args = params.as_array().unwrap();
 
-        let network = &args[0];
-        let token = &args[1];
-        let address = &args[2];
-        let amount = &args[3];
+        let _network = &args[0];
+        let _token = &args[1];
+        let _address = &args[2];
+        let _amount = &args[3];
 
         // 2. Cashier checks if they support the network, and if so,
         //    return adeposit address.
@@ -244,14 +292,13 @@ async fn main() -> Result<()> {
     let (signal, shutdown) = async_channel::unbounded::<()>();
 
     let cashierd2 = cashierd.clone();
-    let cashierd3 = cashierd.clone();
     let (_, _result) = Parallel::new()
         // Run four executor threads.
         .each(0..3, |_| smol::future::block_on(ex.run(shutdown.recv())))
         // Run the main future on the current thread.
         .finish(|| {
             smol::future::block_on(async move {
-                cashierd2.start(ex2, cashierd3.clone().config).await?;
+                cashierd2.start(ex2).await?;
                 drop(signal);
                 Ok::<(), Error>(())
             })