diff --git a/rust/babblerd/src/lib.rs b/rust/babblerd/src/lib.rs index 1dc97981..30bd0bb7 100644 --- a/rust/babblerd/src/lib.rs +++ b/rust/babblerd/src/lib.rs @@ -135,32 +135,23 @@ pub mod if_watcher { } pub mod babel { #[tracing::instrument(skip(sock, receiver))] - pub async fn handle_listener( - sock: UnixStream, - mut receiver: broadcast::Receiver, - ) -> std::io::Result<()> { + pub async fn handle_listener(sock: UnixStream, mut receiver: broadcast::Receiver) { tracing::info!("new socket conn"); let (read, mut write) = sock.into_split(); let mut read = BufReader::new(read).lines(); - let res = loop { + loop { tokio::select! { read = read.next_line() => { - if let None = read? { - break Ok(()); - } + let Ok(Some(_)) = read else { break; }; } recv = receiver.recv() => { - if let Ok(s) = recv { - if let Err(e) = write.write_all(format!("{s}\n").as_bytes()).await { - break Err(e) - }; - continue; - } + let Ok(s) = recv else { break; }; + let Ok(()) = write.write_all(format!("{s}\n").as_bytes()).await else { break; }; } }; - }; + } tracing::info!("closing socket conn"); - write.shutdown().await.and(res) + _ = write.shutdown().await; } #[cfg(target_os = "macos")] @@ -189,6 +180,14 @@ pub mod babel { pub struct BabeldProcess { proc: tokio::process::Child, } + impl Drop for BabeldProcess { + fn drop(&mut self) { + // emergency sigkill babeld process to prevent leakage + if self.proc.try_wait().is_err() { + _ = self.proc.start_kill(); + } + } + } impl BabeldProcess { #[tracing::instrument] fn spawn(first_iface: String) -> Result { @@ -339,7 +338,7 @@ pub mod babel { let mut babel = BabeldProcess::spawn(first_iface)?; let ret = babel.supervise(recv, send).await; - babel.shutdown().await?; + _ = babel.shutdown().await; ret } } diff --git a/rust/babblerd/src/main.rs b/rust/babblerd/src/main.rs index e48fa3ff..6c92ba8e 100644 --- a/rust/babblerd/src/main.rs +++ b/rust/babblerd/src/main.rs @@ -16,7 +16,7 @@ enum State { recv: broadcast::Receiver, babel: JoinHandle>, watcher: JoinHandle>, - listeners: JoinSet>, + listeners: JoinSet<()>, }, } @@ -79,15 +79,15 @@ async fn inner_main() -> color_eyre::Result<()> { sig?; drop(recv); watcher.abort(); - babel.await??; _ = watcher.await; + babel.await??; while let Some(res) = listeners.join_next().await { - res??; + res?; } break; } next_join_result = listeners.join_next(), if !listeners.is_empty() => { - next_join_result.expect("checked")??; + next_join_result.expect("checked")?; tracing::info!("dropped a listener"); if listeners.is_empty() { tracing::info!("stopping babeld"); @@ -111,7 +111,7 @@ async fn inner_main() -> color_eyre::Result<()> { drop(recv); babel.await??; while let Some(res) = listeners.join_next().await { - res??; + res?; } State::Idle } @@ -121,7 +121,7 @@ async fn inner_main() -> color_eyre::Result<()> { watcher.abort(); _ = watcher.await; while let Some(res) = listeners.join_next().await { - res??; + res?; } State::Idle }