From f686d3fe2a2933c65471bfe5b17f0cf2cf48eceb Mon Sep 17 00:00:00 2001 From: Philipp Oppermann Date: Fri, 11 Nov 2022 15:10:30 +0100 Subject: [PATCH] Fix: Spawn control listener threads using `std` API The `tokio::task::spawn_blocking` function will try to join them at the end of `main`, which hangs forever since the listener never finishes. --- binaries/coordinator/src/control.rs | 5 ++--- binaries/coordinator/src/lib.rs | 2 +- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/binaries/coordinator/src/control.rs b/binaries/coordinator/src/control.rs index 7bb21697..a2066952 100644 --- a/binaries/coordinator/src/control.rs +++ b/binaries/coordinator/src/control.rs @@ -9,14 +9,13 @@ use std::{ }; use tokio::sync::{mpsc, oneshot}; use tokio_stream::wrappers::ReceiverStream; -use uuid::Uuid; pub(crate) async fn control_events( control_listen_addr: SocketAddr, ) -> eyre::Result> { let (tx, rx) = mpsc::channel(10); - tokio::task::spawn_blocking(move || listen(control_listen_addr, tx)); + std::thread::spawn(move || listen(control_listen_addr, tx)); Ok(ReceiverStream::new(rx).map(Event::Control)) } @@ -38,7 +37,7 @@ fn listen(control_listen_addr: SocketAddr, tx: mpsc::Sender) { match connection.wrap_err("failed to connect") { Ok(connection) => { let tx = tx.clone(); - tokio::task::spawn_blocking(|| handle_requests(connection, tx)); + std::thread::spawn(|| handle_requests(connection, tx)); } Err(err) => { if tx.blocking_send(err.into()).is_err() { diff --git a/binaries/coordinator/src/lib.rs b/binaries/coordinator/src/lib.rs index 68621bc5..997b85f4 100644 --- a/binaries/coordinator/src/lib.rs +++ b/binaries/coordinator/src/lib.rs @@ -157,7 +157,7 @@ async fn start(runtime_path: &Path) -> eyre::Result<()> { }; let _ = reply_sender.send(reply); } - ControlEvent::Error(_) => todo!(), + ControlEvent::Error(err) => return Err(err), }, } }