Преглед изворни кода

add executor across the whole project to run on more than one thread

ghassmo пре 4 година
родитељ
комит
f3d6156c4c
8 измењених фајлова са 196 додато и 98 уклоњено
  1. 85 37
      src/bin/cashierd.rs
  2. 54 33
      src/bin/darkfid.rs
  3. 14 6
      src/client.rs
  4. 12 6
      src/rpc/rpcserver.rs
  5. 10 2
      src/service/bridge.rs
  6. 5 2
      src/service/btc.rs
  7. 2 2
      src/service/gateway.rs
  8. 14 10
      src/service/sol.rs

+ 85 - 37
src/bin/cashierd.rs

@@ -3,8 +3,10 @@ use std::collections::HashMap;
 use std::path::PathBuf;
 use std::str::FromStr;
 
+use async_executor::Executor;
 use async_trait::async_trait;
 use clap::clap_app;
+use easy_parallel::Parallel;
 use ff::Field;
 use log::debug;
 use rand::rngs::OsRng;
@@ -55,7 +57,7 @@ struct Cashierd {
 
 #[async_trait]
 impl RequestHandler for Cashierd {
-    async fn handle_request(&self, req: JsonRequest) -> JsonResult {
+    async fn handle_request(&self, req: JsonRequest, executor: Arc<Executor<'_>>) -> JsonResult {
         if req.params.as_array().is_none() {
             return JsonResult::Err(jsonerr(InvalidParams, None, req.id));
         }
@@ -63,7 +65,7 @@ impl RequestHandler for Cashierd {
         debug!(target: "RPC", "--> {}", serde_json::to_string(&req).unwrap());
 
         match req.method.as_str() {
-            Some("deposit") => return self.deposit(req.id, req.params).await,
+            Some("deposit") => return self.deposit(req.id, req.params, executor).await,
             Some("withdraw") => return self.withdraw(req.id, req.params).await,
             Some("features") => return self.features(req.id, req.params).await,
             Some(_) => {}
@@ -106,6 +108,7 @@ impl Cashierd {
         bridge: Arc<Bridge>,
         cashier_wallet: Arc<CashierDb>,
         networks: Vec<Network>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<()> {
         debug!(target: "CASHIER DAEMON", "Resume watch deposit keys");
 
@@ -120,6 +123,7 @@ impl Cashierd {
                     .subscribe(
                         deposit_token.drk_public_key,
                         Some(deposit_token.mint_address),
+                        executor.clone(),
                     )
                     .await;
 
@@ -142,6 +146,7 @@ impl Cashierd {
         bridge: Arc<Bridge>,
         cashier_wallet: Arc<CashierDb>,
         recv_coin: async_channel::Receiver<(jubjub::SubgroupPoint, u64)>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<()> {
         // received drk coin
         let (drk_pub_key, amount) = recv_coin.recv().await?;
@@ -155,7 +160,11 @@ impl Cashierd {
         // received drk coin to token publickey
         if let Some(withdraw_token) = token {
             let bridge_subscribtion = bridge
-                .subscribe(drk_pub_key, Some(withdraw_token.mint_address))
+                .subscribe(
+                    drk_pub_key,
+                    Some(withdraw_token.mint_address),
+                    executor.clone(),
+                )
                 .await;
 
             // send a request to the bridge to send amount of token
@@ -199,7 +208,7 @@ impl Cashierd {
         Ok(())
     }
 
-    async fn deposit(&self, id: Value, params: Value) -> JsonResult {
+    async fn deposit(&self, id: Value, params: Value, executor: Arc<Executor<'_>>) -> JsonResult {
         debug!(target: "CASHIER DAEMON", "RECEIVED DEPOSIT REQUEST");
 
         let args: &Vec<serde_json::Value> = params.as_array().unwrap();
@@ -276,7 +285,9 @@ impl Cashierd {
             // record in cashierdb with the network name and token id
 
             let bridge = self.bridge.clone();
-            let bridge_subscribtion = bridge.subscribe(drk_pub_key, mint_address_opt).await;
+            let bridge_subscribtion = bridge
+                .subscribe(drk_pub_key, mint_address_opt, executor)
+                .await;
 
             if check.is_empty() {
                 bridge_subscribtion
@@ -458,6 +469,7 @@ impl Cashierd {
         &mut self,
         mut client: Client,
         state: Arc<Mutex<State>>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<(
         smol::Task<Result<()>>,
         smol::Task<Result<()>>,
@@ -551,10 +563,11 @@ impl Cashierd {
             }
         }
 
-        let resume_watch_deposit_keys_task = smol::spawn(Self::resume_watch_deposit_keys(
+        let resume_watch_deposit_keys_task = executor.spawn(Self::resume_watch_deposit_keys(
             self.bridge.clone(),
             self.cashier_wallet.clone(),
             self.networks.clone(),
+            executor.clone(),
         ));
 
         client.start().await?;
@@ -562,17 +575,25 @@ impl Cashierd {
         let (notify, recv_coin) = async_channel::unbounded::<(jubjub::SubgroupPoint, u64)>();
 
         client
-            .connect_to_subscriber_from_cashier(state, self.cashier_wallet.clone(), notify.clone())
+            .connect_to_subscriber_from_cashier(
+                state,
+                self.cashier_wallet.clone(),
+                notify.clone(),
+                executor.clone(),
+            )
             .await?;
 
         let cashier_wallet = self.cashier_wallet.clone();
         let bridge = self.bridge.clone();
-        let listen_for_receiving_coins_task: smol::Task<Result<()>> = smol::spawn(async move {
+        let ex = executor.clone();
+        let listen_for_receiving_coins_task: smol::Task<Result<()>> = executor.spawn(async move {
+            let ex2 = ex.clone();
             loop {
                 Self::listen_for_receiving_coins(
                     bridge.clone(),
                     cashier_wallet.clone(),
                     recv_coin.clone(),
+                    ex2.clone(),
                 )
                 .await?;
             }
@@ -580,7 +601,7 @@ impl Cashierd {
 
         let bridge2 = self.bridge.clone();
         let listen_for_notification_from_bridge_task: smol::Task<Result<()>> =
-            smol::spawn(async move {
+            executor.spawn(async move {
                 while let Some(token_notification) = bridge2.clone().listen().await {
                     debug!(target: "CASHIER DAEMON", "Notification from birdge");
 
@@ -612,31 +633,11 @@ impl Cashierd {
     }
 }
 
-#[async_std::main]
-async fn main() -> Result<()> {
-    let args = clap_app!(cashierd =>
-        (@arg CONFIG: -c --config +takes_value "Sets a custom config file")
-        (@arg ADDRESS: -a --address "Get Cashier Public key")
-        (@arg verbose: -v --verbose "Increase verbosity")
-    )
-    .get_matches();
-
-    let config_path = if args.is_present("CONFIG") {
-        PathBuf::from(args.value_of("CONFIG").unwrap())
-    } else {
-        join_config_path(&PathBuf::from("cashierd.toml"))?
-    };
-
-    let loglevel = if args.is_present("verbose") {
-        log::Level::Debug
-    } else {
-        log::Level::Info
-    };
-
-    simple_logger::init_with_level(loglevel)?;
-
-    let config: CashierdConfig = Config::<CashierdConfig>::load(config_path)?;
-
+async fn start(
+    executor: Arc<Executor<'_>>,
+    config: &CashierdConfig,
+    get_address_flag: bool,
+) -> Result<()> {
     let mut cashierd = Cashierd::new(config.clone()).await?;
 
     let client_wallet = WalletDb::new(
@@ -693,7 +694,7 @@ async fn main() -> Result<()> {
         public_keys: cashier_public_keys,
     }));
 
-    if args.is_present("ADDRESS") {
+    if get_address_flag {
         let cashier_public = client.main_keypair.public;
         let cashier_public = bs58::encode(&serialize(&cashier_public)).into_string();
         println!("Public Key: {}", cashier_public);
@@ -707,8 +708,8 @@ async fn main() -> Result<()> {
         identity_pass: config.tls_identity_password.clone(),
     };
 
-    let (t1, t2, t3) = cashierd.start(client, state).await?;
-    listen_and_serve(cfg, Arc::new(cashierd)).await?;
+    let (t1, t2, t3) = cashierd.start(client, state, executor.clone()).await?;
+    listen_and_serve(cfg, Arc::new(cashierd), executor).await?;
 
     t1.cancel().await;
     t2.cancel().await;
@@ -716,3 +717,50 @@ async fn main() -> Result<()> {
 
     Ok(())
 }
+
+#[async_std::main]
+async fn main() -> Result<()> {
+    let args = clap_app!(cashierd =>
+        (@arg CONFIG: -c --config +takes_value "Sets a custom config file")
+        (@arg ADDRESS: -a --address "Get Cashier Public key")
+        (@arg verbose: -v --verbose "Increase verbosity")
+    )
+    .get_matches();
+
+    let config_path = if args.is_present("CONFIG") {
+        PathBuf::from(args.value_of("CONFIG").unwrap())
+    } else {
+        join_config_path(&PathBuf::from("cashierd.toml"))?
+    };
+
+    let loglevel = if args.is_present("verbose") {
+        log::Level::Debug
+    } else {
+        log::Level::Info
+    };
+
+    simple_logger::init_with_level(loglevel)?;
+
+    let config: CashierdConfig = Config::<CashierdConfig>::load(config_path)?;
+
+    let ex = Arc::new(Executor::new());
+    let (signal, shutdown) = async_channel::unbounded::<()>();
+
+    let ex2 = ex.clone();
+
+    let get_address_flag = args.is_present("ADDRESS");
+
+    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 {
+                start(ex2, &config, get_address_flag).await?;
+                drop(signal);
+                Ok::<(), drk::Error>(())
+            })
+        });
+
+    result
+}

+ 54 - 33
src/bin/darkfid.rs

@@ -3,8 +3,10 @@ use std::collections::HashMap;
 use std::path::PathBuf;
 use std::str::FromStr;
 
+use async_executor::Executor;
 use async_trait::async_trait;
 use clap::clap_app;
+use easy_parallel::Parallel;
 use log::debug;
 use serde_json::{json, Value};
 
@@ -38,7 +40,7 @@ pub struct Cashier {
 
 #[async_trait]
 impl RequestHandler for Darkfid {
-    async fn handle_request(&self, req: JsonRequest) -> JsonResult {
+    async fn handle_request(&self, req: JsonRequest, _executor: Arc<Executor<'_>>) -> JsonResult {
         if req.params.as_array().is_none() {
             return JsonResult::Err(jsonerr(InvalidParams, None, req.id));
         }
@@ -81,12 +83,12 @@ impl Darkfid {
         })
     }
 
-    async fn start(&mut self, state: Arc<Mutex<State>>) -> Result<()> {
+    async fn start(&mut self, state: Arc<Mutex<State>>, executor: Arc<Executor<'_>>) -> Result<()> {
         self.client.lock().await.start().await?;
         self.client
             .lock()
             .await
-            .connect_to_subscriber(state)
+            .connect_to_subscriber(state, executor)
             .await?;
 
         Ok(())
@@ -136,10 +138,8 @@ impl Darkfid {
             let mut symbols: HashMap<String, (String, String)> = HashMap::new();
 
             for balance in balances.list.iter() {
-
-                
-                // XXX: this must be changed once cashierd 
-                // supports more than two networks 
+                // XXX: this must be changed once cashierd
+                // supports more than two networks
 
                 let mut network = "solana";
 
@@ -496,30 +496,7 @@ impl Darkfid {
     }
 }
 
-#[async_std::main]
-async fn main() -> Result<()> {
-    let args = clap_app!(darkfid =>
-        (@arg CONFIG: -c --config +takes_value "Sets a custom config file")
-        (@arg verbose: -v --verbose "Increase verbosity")
-    )
-    .get_matches();
-
-    let config_path = if args.is_present("CONFIG") {
-        PathBuf::from(args.value_of("CONFIG").unwrap())
-    } else {
-        join_config_path(&PathBuf::from("darkfid.toml"))?
-    };
-
-    let loglevel = if args.is_present("verbose") {
-        log::Level::Debug
-    } else {
-        log::Level::Info
-    };
-
-    simple_logger::init_with_level(loglevel)?;
-
-    let config: DarkfidConfig = Config::<DarkfidConfig>::load(config_path)?;
-
+async fn start(executor: Arc<Executor<'_>>, config: &DarkfidConfig) -> Result<()> {
     let wallet = WalletDb::new(
         expand_path(&config.wallet_path)?.as_path(),
         config.wallet_password.clone(),
@@ -601,6 +578,50 @@ async fn main() -> Result<()> {
         identity_pass: config.tls_identity_password.clone(),
     };
 
-    darkfid.start(state).await?;
-    listen_and_serve(server_config, Arc::new(darkfid)).await
+    darkfid.start(state, executor.clone()).await?;
+    listen_and_serve(server_config, Arc::new(darkfid), executor).await
+}
+
+#[async_std::main]
+async fn main() -> Result<()> {
+    let args = clap_app!(darkfid =>
+        (@arg CONFIG: -c --config +takes_value "Sets a custom config file")
+        (@arg verbose: -v --verbose "Increase verbosity")
+    )
+    .get_matches();
+
+    let config_path = if args.is_present("CONFIG") {
+        PathBuf::from(args.value_of("CONFIG").unwrap())
+    } else {
+        join_config_path(&PathBuf::from("darkfid.toml"))?
+    };
+
+    let loglevel = if args.is_present("verbose") {
+        log::Level::Debug
+    } else {
+        log::Level::Info
+    };
+
+    simple_logger::init_with_level(loglevel)?;
+
+    let config: DarkfidConfig = Config::<DarkfidConfig>::load(config_path)?;
+
+    let ex = Arc::new(Executor::new());
+    let (signal, shutdown) = async_channel::unbounded::<()>();
+
+    let ex2 = ex.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 {
+                start(ex2, &config).await?;
+                drop(signal);
+                Ok::<(), drk::Error>(())
+            })
+        });
+
+    result
 }

+ 14 - 6
src/client.rs

@@ -1,5 +1,6 @@
 use std::net::SocketAddr;
 
+use async_executor::Executor;
 use async_std::sync::{Arc, Mutex};
 use bellman::groth16;
 use bls12_381::Bls12;
@@ -18,7 +19,7 @@ use crate::{
     service::{GatewayClient, GatewaySlabsSubscriber},
     state::{state_transition, ProgramState, StateUpdate},
     tx,
-    wallet::{CashierDbPtr, Keypair, WalletPtr, walletdb::Balances},
+    wallet::{walletdb::Balances, CashierDbPtr, Keypair, WalletPtr},
     Result,
 };
 
@@ -235,15 +236,17 @@ impl Client {
         state: Arc<Mutex<State>>,
         cashier_wallet: CashierDbPtr,
         notify: async_channel::Sender<(jubjub::SubgroupPoint, u64)>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<()> {
         // start subscribing
         debug!(target: "CLIENT", "Start subscriber for cashier");
-        let gateway_slabs_sub: GatewaySlabsSubscriber = self.gateway.start_subscriber().await?;
+        let gateway_slabs_sub: GatewaySlabsSubscriber =
+            self.gateway.start_subscriber(executor.clone()).await?;
 
         let secret_key = self.main_keypair.private;
         let wallet = self.wallet.clone();
 
-        let task: smol::Task<Result<()>> = smol::spawn(async move {
+        let task: smol::Task<Result<()>> = executor.spawn(async move {
             loop {
                 let slab = gateway_slabs_sub.recv().await?;
 
@@ -291,15 +294,20 @@ impl Client {
         Ok(())
     }
 
-    pub async fn connect_to_subscriber(&self, state: Arc<Mutex<State>>) -> Result<()> {
+    pub async fn connect_to_subscriber(
+        &self,
+        state: Arc<Mutex<State>>,
+        executor: Arc<Executor<'_>>,
+    ) -> Result<()> {
         // start subscribing
         debug!(target: "CLIENT", "Start subscriber");
-        let gateway_slabs_sub: GatewaySlabsSubscriber = self.gateway.start_subscriber().await?;
+        let gateway_slabs_sub: GatewaySlabsSubscriber =
+            self.gateway.start_subscriber(executor.clone()).await?;
 
         let secret_key = self.main_keypair.private;
         let wallet = self.wallet.clone();
 
-        let task: smol::Task<Result<()>> = smol::spawn(async move {
+        let task: smol::Task<Result<()>> = executor.spawn(async move {
             loop {
                 let slab = gateway_slabs_sub.recv().await?;
 

+ 12 - 6
src/rpc/rpcserver.rs

@@ -2,6 +2,7 @@ use std::net::{SocketAddr, TcpListener, TcpStream};
 use std::path::PathBuf;
 use std::sync::Arc;
 
+use async_executor::Executor;
 use async_native_tls::{Identity, TlsAcceptor};
 use async_trait::async_trait;
 use log::{debug, error};
@@ -22,13 +23,14 @@ pub struct RpcServerConfig {
 
 #[async_trait]
 pub trait RequestHandler: Sync + Send {
-    async fn handle_request(&self, req: JsonRequest) -> JsonResult;
+    async fn handle_request(&self, req: JsonRequest, executor: Arc<Executor<'_>>) -> JsonResult;
 }
 
 async fn serve(
     mut stream: Async<TcpStream>,
     tls: Option<TlsAcceptor>,
     rh: Arc<impl RequestHandler + 'static>,
+    executor: Arc<Executor<'_>>,
 ) -> Result<()> {
     debug!(target: "RPC SERVER", "Accepted connection");
 
@@ -58,7 +60,7 @@ async fn serve(
                 }
             };
 
-            let reply = rh.handle_request(r).await;
+            let reply = rh.handle_request(r, executor.clone()).await;
             let j = serde_json::to_string(&reply).unwrap();
             debug!(target: "RPC", "<-- {}", j);
 
@@ -92,7 +94,7 @@ async fn serve(
                     }
                 };
 
-                let reply = rh.handle_request(r).await;
+                let reply = rh.handle_request(r, executor.clone()).await;
                 let j = serde_json::to_string(&reply).unwrap();
                 debug!(target: "RPC", "<-- {}", j);
 
@@ -113,6 +115,7 @@ async fn listen(
     listener: Async<TcpListener>,
     tls: Option<TlsAcceptor>,
     rh: Arc<impl RequestHandler + 'static>,
+    executor: Arc<Executor<'_>>,
 ) -> Result<()> {
     match &tls {
         None => {
@@ -123,13 +126,15 @@ async fn listen(
         }
     }
 
+    let ex = executor.clone();
     loop {
         let (stream, _) = listener.accept().await?;
         let tls = tls.clone();
         let rh_c = rh.clone();
 
-        smol::spawn(async move {
-            if let Err(err) = serve(stream, tls, rh_c).await {
+        let ex2 = ex.clone();
+        ex.spawn(async move {
+            if let Err(err) = serve(stream, tls, rh_c, ex2.clone()).await {
                 error!(target: "RPC SERVER", "Connection error: {:#?}", err);
             }
         })
@@ -140,6 +145,7 @@ async fn listen(
 pub async fn listen_and_serve(
     cfg: RpcServerConfig,
     rh: Arc<impl RequestHandler + 'static>,
+    executor: Arc<Executor<'_>>,
 ) -> Result<()> {
     let tls: Option<TlsAcceptor>;
 
@@ -151,6 +157,6 @@ pub async fn listen_and_serve(
         tls = None;
     }
 
-    let listener = listen(Async::<TcpListener>::bind(cfg.socket_addr)?, tls, rh);
+    let listener = listen(Async::<TcpListener>::bind(cfg.socket_addr)?, tls, rh, executor);
     listener.await
 }

+ 10 - 2
src/service/bridge.rs

@@ -1,6 +1,7 @@
 use async_std::sync::{Arc, Mutex};
 use std::collections::HashMap;
 
+use async_executor::Executor;
 use async_trait::async_trait;
 use futures::stream::FuturesUnordered;
 use futures::stream::StreamExt;
@@ -116,12 +117,15 @@ impl Bridge {
         self: Arc<Self>,
         drk_pub_key: jubjub::SubgroupPoint,
         mint: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> BridgeSubscribtion {
         debug!(target: "BRIDGE", "Start new subscription");
         let (sender, req) = async_channel::unbounded();
         let (rep, receiver) = async_channel::unbounded();
 
-        smol::spawn(self.listen_for_new_subscription(req, rep, drk_pub_key, mint)).detach();
+        executor
+            .spawn(self.listen_for_new_subscription(req, rep, drk_pub_key, mint, executor.clone()))
+            .detach();
 
         BridgeSubscribtion { sender, receiver }
     }
@@ -132,6 +136,7 @@ impl Bridge {
         rep: async_channel::Sender<BridgeResponse>,
         drk_pub_key: jubjub::SubgroupPoint,
         mint: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<()> {
         debug!(target: "BRIDGE", "Listen for new subscription");
         let req = req.recv().await?;
@@ -171,6 +176,7 @@ impl Bridge {
                             token_key.public_key,
                             drk_pub_key,
                             mint_address,
+                            executor
                         )
                         .await;
 
@@ -188,7 +194,7 @@ impl Bridge {
                     }
                 }
                 None => {
-                    let sub = client.subscribe(drk_pub_key, mint_address).await;
+                    let sub = client.subscribe(drk_pub_key, mint_address, executor).await;
                     if sub.is_err() {
                         error!(target: "BRIDGE", "{}", sub.unwrap_err().to_string());
                         res = BridgeResponse {
@@ -234,6 +240,7 @@ pub trait NetworkClient {
         self: Arc<Self>,
         drk_pub_key: jubjub::SubgroupPoint,
         mint: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<TokenSubscribtion>;
 
     // should check if the keypair in not already subscribed
@@ -243,6 +250,7 @@ pub trait NetworkClient {
         public_key: Vec<u8>,
         drk_pub_key: jubjub::SubgroupPoint,
         mint: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<String>;
 
     async fn get_notifier(self: Arc<Self>) -> Result<async_channel::Receiver<TokenNotification>>;

+ 5 - 2
src/service/btc.rs

@@ -3,6 +3,7 @@ use std::convert::From;
 use std::str::FromStr;
 use std::time::Duration;
 
+use async_executor::Executor;
 use async_trait::async_trait;
 
 use bitcoin::blockdata::{
@@ -348,6 +349,7 @@ impl NetworkClient for BtcClient {
         self: Arc<Self>,
         drk_pub_key: jubjub::SubgroupPoint,
         _mint: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<TokenSubscribtion> {
         // Generate bitcoin keys
         let keypair = Keypair::new();
@@ -358,7 +360,7 @@ impl NetworkClient for BtcClient {
         // start scheduler for checking balance
         debug!(target: "BRIDGE BITCOIN", "Subscribing for deposit");
 
-        smol::spawn(async move {
+        executor.spawn(async move {
             let result = self.handle_subscribe_request(btc_keys, drk_pub_key).await;
             if let Err(e) = result {
                 error!(target: "BTC BRIDGE SUBSCRIPTION","{}", e.to_string());
@@ -378,12 +380,13 @@ impl NetworkClient for BtcClient {
         _public_key: Vec<u8>,
         drk_pub_key: jubjub::SubgroupPoint,
         _mint: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<String> {
         let keypair: Keypair = deserialize(&private_key)?;
         let btc_keys = Account::new(&keypair, self.network);
         let public_key = keypair.pubkey().to_string();
 
-        smol::spawn(async move {
+        executor.spawn(async move {
             let result = self.handle_subscribe_request(btc_keys, drk_pub_key).await;
             if let Err(e) = result {
                 error!(target: "BTC BRIDGE SUBSCRIPTION","{}", e.to_string());

+ 2 - 2
src/service/gateway.rs

@@ -290,12 +290,12 @@ impl GatewayClient {
         self.slabstore.clone()
     }
 
-    pub async fn start_subscriber(&self) -> Result<GatewaySlabsSubscriber> {
+    pub async fn start_subscriber(&self, executor: Arc<Executor<'_>>) -> Result<GatewaySlabsSubscriber> {
         debug!(target: "GATEWAY CLIENT","Start subscriber");
 
         let mut subscriber = Subscriber::new(self.sub_addr, String::from("GATEWAY CLIENT"));
         subscriber.start().await?;
-        smol::spawn(Self::subscribe_loop(
+        executor.spawn(Self::subscribe_loop(
             subscriber,
             self.slabstore.clone(),
             self.gateway_slabs_sub_s.clone(),

+ 14 - 10
src/service/sol.rs

@@ -3,6 +3,7 @@ use std::convert::TryFrom;
 use std::str::FromStr;
 use std::time::Duration;
 
+use async_executor::Executor;
 use async_native_tls::TlsConnector;
 use async_trait::async_trait;
 use futures::{SinkExt, StreamExt};
@@ -403,6 +404,7 @@ impl NetworkClient for SolClient {
         self: Arc<Self>,
         drk_pub_key: jubjub::SubgroupPoint,
         mint_address: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<TokenSubscribtion> {
         let keypair = Keypair::generate(&mut OsRng);
 
@@ -418,7 +420,7 @@ impl NetworkClient for SolClient {
             return Err(Error::from(SolFailed::MainAccountNotEnoughValue));
         }
 
-        smol::spawn(async move {
+        executor.spawn(async move {
             let result = self
                 .handle_subscribe_request(keypair, drk_pub_key, mint)
                 .await;
@@ -441,6 +443,7 @@ impl NetworkClient for SolClient {
         _public_key: Vec<u8>,
         drk_pub_key: jubjub::SubgroupPoint,
         mint_address: Option<String>,
+        executor: Arc<Executor<'_>>,
     ) -> Result<String> {
         let keypair: Keypair = deserialize(&private_key)?;
 
@@ -454,15 +457,16 @@ impl NetworkClient for SolClient {
             return Err(Error::from(SolFailed::MainAccountNotEnoughValue));
         }
 
-        smol::spawn(async move {
-            let result = self
-                .handle_subscribe_request(keypair, drk_pub_key, mint)
-                .await;
-            if let Err(e) = result {
-                error!(target: "SOL BRIDGE SUBSCRIPTION","{}", e.to_string());
-            }
-        })
-        .detach();
+        executor
+            .spawn(async move {
+                let result = self
+                    .handle_subscribe_request(keypair, drk_pub_key, mint)
+                    .await;
+                if let Err(e) = result {
+                    error!(target: "SOL BRIDGE SUBSCRIPTION","{}", e.to_string());
+                }
+            })
+            .detach();
 
         Ok(public_key)
     }