Browse Source

rpc/server: detach request handling

aggstam 3 years ago
parent
commit
3692735012

+ 2 - 1
bin/darkfid/src/main.rs

@@ -419,7 +419,8 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
 
     // JSON-RPC server
     info!("Starting JSON-RPC server");
-    ex.spawn(listen_and_serve(args.rpc_listen, darkfid.clone())).detach();
+    let _ex = ex.clone();
+    ex.spawn(listen_and_serve(args.rpc_listen, darkfid.clone(), _ex)).detach();
 
     info!("Starting sync P2P network");
     sync_p2p.clone().unwrap().start(ex.clone()).await?;

+ 2 - 1
bin/darkwiki/darkwikid/src/main.rs

@@ -598,7 +598,8 @@ async fn realmain(args: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
     // JSON-RPC server
     // ===============
     let rpc_iface = Arc::new(JsonRpcInterface::new(rpc_tx, notify_rx));
-    executor.spawn(listen_and_serve(args.rpc_listen, rpc_iface)).detach();
+    let _ex = executor.clone();
+    executor.spawn(listen_and_serve(args.rpc_listen, rpc_iface, _ex)).detach();
 
     // ====
     // Raft

+ 2 - 1
bin/faucetd/src/main.rs

@@ -611,7 +611,8 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
 
     // JSON-RPC server
     info!("Starting JSON-RPC server");
-    ex.spawn(listen_and_serve(args.rpc_listen, faucetd.clone())).detach();
+    let _ex = ex.clone();
+    ex.spawn(listen_and_serve(args.rpc_listen, faucetd.clone(), _ex)).detach();
 
     info!("Starting sync P2P network");
     sync_p2p.clone().start(ex.clone()).await?;

+ 2 - 1
bin/fud/fud/src/main.rs

@@ -420,7 +420,8 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
 
     // JSON-RPC server
     info!("Starting JSON-RPC server");
-    ex.spawn(listen_and_serve(args.rpc_listen, fud.clone())).detach();
+    let _ex = ex.clone();
+    ex.spawn(listen_and_serve(args.rpc_listen, fud.clone(), _ex)).detach();
 
     info!("Starting sync P2P network");
     p2p.clone().start(ex.clone()).await?;

+ 4 - 1
bin/ircd/src/main.rs

@@ -168,7 +168,10 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
     let rpc_listen_addr = settings.rpc_listen.clone();
     let rpc_interface =
         Arc::new(JsonRpcInterface { addr: rpc_listen_addr.clone(), p2p: p2p.clone() });
-    executor.spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface).await }).detach();
+    let _ex = executor.clone();
+    executor
+        .spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface, _ex).await })
+        .detach();
 
     //
     // IRC instance

+ 4 - 1
bin/ircd2/src/main.rs

@@ -123,7 +123,10 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
     let rpc_listen_addr = settings.rpc_listen.clone();
     let rpc_interface =
         Arc::new(JsonRpcInterface { addr: rpc_listen_addr.clone(), p2p: p2p.clone() });
-    executor.spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface).await }).detach();
+    let _ex = executor.clone();
+    executor
+        .spawn(async move { listen_and_serve(rpc_listen_addr, rpc_interface, _ex).await })
+        .detach();
 
     ////////////////////
     // IRC server

+ 2 - 2
bin/lilith/src/main.rs

@@ -342,10 +342,10 @@ async fn realmain(args: Args, ex: Arc<smol::Executor<'_>>) -> Result<()> {
 
     // JSON-RPC server
     info!("Starting JSON-RPC server");
-    ex.spawn(listen_and_serve(args.rpc_listen, lilith.clone())).detach();
+    let _ex = ex.clone();
+    ex.spawn(listen_and_serve(args.rpc_listen, lilith.clone(), _ex)).detach();
 
     // JSON-RPC notifications simulation
-    let _ex = ex.clone();
     ex.spawn(simulate_blocks(subscriber)).detach();
 
     // Wait for SIGINT

+ 2 - 1
bin/tau/taud/src/main.rs

@@ -276,7 +276,8 @@ async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<(
         workspaces.clone(),
         p2p.clone(),
     ));
-    executor.spawn(listen_and_serve(settings.rpc_listen.clone(), rpc_interface)).detach();
+    let _ex = executor.clone();
+    executor.spawn(listen_and_serve(settings.rpc_listen.clone(), rpc_interface, _ex)).detach();
 
     //
     // Waiting Exit signal

+ 2 - 1
example/dchat/src/main.rs

@@ -246,7 +246,8 @@ async fn main() -> Result<()> {
     // ANCHOR: json_init
     let accept_addr = settings.accept_addr.clone();
     let rpc = Arc::new(JsonRpcInterface { addr: accept_addr.clone(), p2p });
-    ex.spawn(async move { listen_and_serve(accept_addr.clone(), rpc).await }).detach();
+    let _ex = ex.clone();
+    ex.spawn(async move { listen_and_serve(accept_addr.clone(), rpc, _ex).await }).detach();
     // ANCHOR_END: json_init
 
     let nthreads = num_cpus::get();

+ 2 - 1
example/p2pdebug/src/main.rs

@@ -246,7 +246,8 @@ impl MockP2p {
 
         let rpc_interface =
             Arc::new(rpc::JsonRpcInterface { addr: rpc_addr.clone(), p2p: p2p.clone() });
-        executor.spawn(async move { listen_and_serve(rpc_addr, rpc_interface).await }).detach();
+        let _ex = executor.clone();
+        executor.spawn(async move { listen_and_serve(rpc_addr, rpc_interface, _ex).await }).detach();
 
         p2p.clone().start(executor.clone()).await?;
         p2p.run(executor).await

+ 2 - 1
script/research/dhtd/src/main.rs

@@ -299,7 +299,8 @@ async fn realmain(args: Args, ex: Arc<Executor<'_>>) -> Result<()> {
 
     // JSON-RPC server
     info!("Starting JSON-RPC server");
-    ex.spawn(listen_and_serve(args.rpc_listen, dhtd.clone())).detach();
+    let _ex = ex.clone();
+    ex.spawn(listen_and_serve(args.rpc_listen, dhtd.clone(), _ex)).detach();
 
     info!("Starting sync P2P network");
     p2p.clone().start(ex.clone()).await?;

+ 12 - 4
src/rpc/server.rs

@@ -117,10 +117,17 @@ async fn accept(
 async fn run_accept_loop(
     listener: Box<dyn TransportListener>,
     rh: Arc<impl RequestHandler + 'static>,
+    ex: Arc<smol::Executor<'_>>,
 ) -> Result<()> {
     while let Ok((stream, peer_addr)) = listener.next().await {
         info!("JSON-RPC server accepted connection from {}", peer_addr);
-        accept(stream, peer_addr, rh.clone()).await?;
+        // Detaching requests handling
+        let _rh = rh.clone();
+        ex.spawn(async move {
+            if let Err(e) = accept(stream, peer_addr.clone(), _rh).await {
+                error!(target: "jsonrpc-server", "JSON-RPC server error on handling request of {}: {}", peer_addr, e);
+            }
+        }).detach();
     }
 
     Ok(())
@@ -131,6 +138,7 @@ async fn run_accept_loop(
 pub async fn listen_and_serve(
     accept_url: Url,
     rh: Arc<impl RequestHandler + 'static>,
+    ex: Arc<smol::Executor<'_>>,
 ) -> Result<()> {
     debug!(target: "jsonrpc-server", "Trying to bind listener on {}", accept_url);
 
@@ -151,12 +159,12 @@ pub async fn listen_and_serve(
             match $upgrade {
                 None => {
                     info!("JSON-RPC listener bound to {}", accept_url);
-                    run_accept_loop(Box::new(listener), rh).await?;
+                    run_accept_loop(Box::new(listener), rh, ex.clone()).await?;
                 }
                 Some(u) if u == "tls" => {
                     let tls_listener = $transport.upgrade_listener(listener)?.await?;
                     info!("JSON-RPC listener bound to {}", accept_url);
-                    run_accept_loop(Box::new(tls_listener), rh).await?;
+                    run_accept_loop(Box::new(tls_listener), rh, ex.clone()).await?;
                 }
                 Some(u) => return Err(Error::UnsupportedTransportUpgrade(u)),
             }
@@ -189,7 +197,7 @@ pub async fn listen_and_serve(
                 error!("JSON-RPC Unix socket bind to {} failed: {}", accept_url, err);
                 return Err(Error::BindFailed(accept_url.as_str().into()))
             }
-            run_accept_loop(Box::new(listener?), rh).await?;
+            run_accept_loop(Box::new(listener?), rh, ex.clone()).await?;
         }
         _ => unimplemented!(),
     }