Bladeren bron

create run_rpc_server function and make the rpc server run inside start
function in cashierd

ghassmo 4 jaren geleden
bovenliggende
commit
be457dd4a6
1 gewijzigde bestanden met toevoegingen van 67 en 63 verwijderingen
  1. 67 63
      src/bin/cashierd.rs

+ 67 - 63
src/bin/cashierd.rs

@@ -119,23 +119,24 @@ impl Cashierd {
         let cashier_wallet = self.cashier_wallet.clone();
 
         let ex = executor.clone();
-        executor
-            .spawn(async move {
-                loop {
-                    Self::listen_for_receiving_coins(
-                        ex.clone(),
-                        bridge.clone(),
-                        cashier_wallet.clone(),
-                        recv_coin.clone(),
-                    )
-                    .await
-                    .expect(" listen for receiving coins");
-                }
-            })
-            .await;
+        let listen_for_receiving_coins_task = executor.spawn(async move {
+            loop {
+                Self::listen_for_receiving_coins(
+                    ex.clone(),
+                    bridge.clone(),
+                    cashier_wallet.clone(),
+                    recv_coin.clone(),
+                )
+                .await
+                .expect(" listen for receiving coins");
+            }
+        });
 
-        cashier_client_subscriber_task.cancel().await;
+        let rpc_url = self.config.rpc_url.clone();
+        run_rpc_server(self.clone(), rpc_url).await?;
 
+        listen_for_receiving_coins_task.cancel().await;
+        cashier_client_subscriber_task.cancel().await;
         Ok(())
     }
 
@@ -273,6 +274,56 @@ impl Cashierd {
     }
 }
 
+async fn run_rpc_server(cashierd: Cashierd, rpc_url: String) -> Result<()> {
+    let listener = TcpListener::bind(rpc_url.clone()).await?;
+    debug!(target: "RPC SERVER", "Listening on {}", rpc_url);
+    loop {
+        debug!(target: "RPC SERVER", "waiting for client");
+
+        let (mut socket, _) = listener.accept().await?;
+
+        debug!(target: "RPC SERVER", "accepted client");
+
+        let cashierd = cashierd.clone();
+        tokio::spawn(async move {
+            let mut buf = [0; 2048];
+
+            loop {
+                let n = match socket.read(&mut buf).await {
+                    Ok(n) if n == 0 => {
+                        debug!(target: "RPC SERVER", "closed connection");
+                        return;
+                    }
+                    Ok(n) => n,
+                    Err(e) => {
+                        debug!(target: "RPC SERVER", "failed to read from socket; err = {:?}", e);
+                        return;
+                    }
+                };
+
+                let r: JsonRequest = match serde_json::from_slice(&buf[0..n]) {
+                    Ok(r) => r,
+                    Err(e) => {
+                        debug!(target: "RPC SERVER", "received invalid json; err = {:?}", e);
+                        return;
+                    }
+                };
+
+                let reply = cashierd.clone().handle_request(r).await;
+                let j = serde_json::to_string(&reply).unwrap();
+
+                debug!(target: "RPC", "<-- {:#?}", j);
+
+                // Write the data back
+                if let Err(e) = socket.write_all(j.as_bytes()).await {
+                    debug!(target: "RPC SERVER", "failed to write to socket; err = {:?}", e);
+                    return;
+                }
+            }
+        });
+    }
+}
+
 #[tokio::main]
 async fn main() -> Result<()> {
     let args = clap_app!(cashierd =>
@@ -291,9 +342,6 @@ async fn main() -> Result<()> {
 
     let cashierd = Cashierd::new(args.clone().is_present("verbose"), config_path)?;
 
-    let listener = TcpListener::bind(cashierd.clone().config.rpc_url).await?;
-    debug!(target: "RPC SERVER", "Listening on {}", cashierd.clone().config.rpc_url);
-
     let logger_config = ConfigBuilder::new().set_time_format_str("%T%.6f").build();
     let debug_level = if args.is_present("verbose") {
         LevelFilter::Debug
@@ -329,49 +377,5 @@ async fn main() -> Result<()> {
             })
         });
 
-    loop {
-        debug!(target: "RPC SERVER", "waiting for client");
-
-        let (mut socket, _) = listener.accept().await?;
-
-        debug!(target: "RPC SERVER", "accepted client");
-
-        let cashierd = cashierd.clone();
-        tokio::spawn(async move {
-            let mut buf = [0; 2048];
-
-            loop {
-                let n = match socket.read(&mut buf).await {
-                    Ok(n) if n == 0 => {
-                        debug!(target: "RPC SERVER", "closed connection");
-                        return;
-                    }
-                    Ok(n) => n,
-                    Err(e) => {
-                        debug!(target: "RPC SERVER", "failed to read from socket; err = {:?}", e);
-                        return;
-                    }
-                };
-
-                let r: JsonRequest = match serde_json::from_slice(&buf[0..n]) {
-                    Ok(r) => r,
-                    Err(e) => {
-                        debug!(target: "RPC SERVER", "received invalid json; err = {:?}", e);
-                        return;
-                    }
-                };
-
-                let reply = cashierd.clone().handle_request(r).await;
-                let j = serde_json::to_string(&reply).unwrap();
-
-                debug!(target: "RPC", "<-- {:#?}", j);
-
-                // Write the data back
-                if let Err(e) = socket.write_all(j.as_bytes()).await {
-                    debug!(target: "RPC SERVER", "failed to write to socket; err = {:?}", e);
-                    return;
-                }
-            }
-        });
-    }
+    Ok(())
 }