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), }, } }