From 0b3424857ec1c5b57d6778f2bf6f379685e7f073 Mon Sep 17 00:00:00 2001 From: Lorenz Leutgeb Date: Mon, 13 Oct 2025 17:42:37 +0200 Subject: [PATCH] node/reactor: Correctly handle error events `fn handle_events` would panic, if there were multiple events for one token, and the first one that happened to be handled was an error. Indeed it is concerning if a token is encountered that was never registered before. However, tokens that were just deregistered must be tracked. Using `Vec` here seems a bit costly, in the future, `smallvec::SmallVec` could be considered. The "unregister" methods are renamed to "deregister" to better line up with `mio` vocabulary. Log statements that helped analysis of the panic that occurred here are overhauled and improved, requiring a `Debug` bound on types that obviously implement it. --- crates/radicle-node/src/reactor.rs | 82 ++++++++++++---------- crates/radicle-node/src/reactor/session.rs | 2 +- crates/radicle-node/src/runtime.rs | 7 +- crates/radicle-node/src/test/node.rs | 3 +- crates/radicle-node/src/wire.rs | 3 +- 5 files changed, 57 insertions(+), 40 deletions(-) diff --git a/crates/radicle-node/src/reactor.rs b/crates/radicle-node/src/reactor.rs index 0aeda592..1e0802df 100644 --- a/crates/radicle-node/src/reactor.rs +++ b/crates/radicle-node/src/reactor.rs @@ -203,13 +203,13 @@ pub trait ReactionHandler: Send + Iterator Runtime { /// /// Whether one of the events was originated from the waker. fn handle_events(&mut self, time: LocalTime, events: Events) -> bool { + log::trace!(target: "reactor", "Handling events"); let mut awoken = false; + let mut deregistered = Vec::new(); for event in events.into_iter() { - let id = event.token(); + let token = event.token(); - if id == WAKER { + if token == WAKER { log::trace!(target: "reactor", "Awoken by the controller"); awoken = true; - } else if self.listeners.contains_key(&id) { - log::trace!(target: "reactor", event:debug; "From listener"); + } else if self.listeners.contains_key(&token) { + log::trace!(target: "reactor", token=token.0; "Event from listener with token {}: {:?}", token.0, event); if !event.is_error() { - let listener = self.listeners.get_mut(&id).expect("resource disappeared"); + let listener = self + .listeners + .get_mut(&token) + .expect("resource disappeared"); listener .handle(event) .into_iter() .for_each(|service_event| { - self.service.listener_reacted(id, service_event, time); + self.service.listener_reacted(token, service_event, time); }); } else { - let listener = self - .unregister_listener(id) - .expect("listener has disappeared"); + let listener = self.deregister_listener(token).unwrap_or_else(|| { + panic!("listener with token {} has disappeared", token.0) + }); self.service - .handle_error(Error::ListenerDisconnect(id, listener)); + .handle_error(Error::ListenerDisconnect(token, listener)); + deregistered.push(token); } - } else if self.transports.contains_key(&id) { - log::trace!(target: "reactor", event:debug; "From transport"); + } else if self.transports.contains_key(&token) { + log::trace!(target: "reactor", token=token.0; "Event from transport with token {}: {:?}", token.0, event); if !event.is_error() { - let transport = self.transports.get_mut(&id).expect("resource disappeared"); + let transport = self + .transports + .get_mut(&token) + .expect("resource disappeared"); transport .handle(event) .into_iter() .for_each(|service_event| { - self.service.transport_reacted(id, service_event, time); + self.service.transport_reacted(token, service_event, time); }); } else { - let transport = self - .unregister_transport(id) - .expect("transport has disappeared"); + let transport = self.deregister_transport(token).unwrap_or_else(|| { + panic!("transport with token {} has disappeared", token.0) + }); self.service - .handle_error(Error::TransportDisconnect(id, transport)); + .handle_error(Error::TransportDisconnect(token, transport)); + deregistered.push(token); } - } else { - panic!("token in poll which is not a known waker, listener or transport") + } else if !deregistered.contains(&token) { + log::warn!(target: "reactor", token=token.0; "Event from unknown token {}: {:?}", token.0, event); } } @@ -498,7 +508,7 @@ impl Runtime { ) -> Result<(), Error> { match action { Action::RegisterListener(token, mut listener) => { - log::debug!(target: "reactor", token=token.0; "Registering listener"); + log::trace!(target: "reactor", token=token.0; "Registering listener {:?} with token {}", listener, token.0); self.poll .registry() @@ -520,33 +530,33 @@ impl Runtime { .transport_registered(token, &self.transports[&token]); } Action::UnregisterListener(token) => { - let Some(listener) = self.unregister_listener(token) else { + let Some(listener) = self.deregister_listener(token) else { return Ok(()); }; - log::debug!(target: "reactor", token=token.0; "Handing over listener"); + log::debug!(target: "reactor", token=token.0; "Handing over listener {listener:?} with token {}", token.0); self.service.handover_listener(token, listener); } Action::UnregisterTransport(token) => { - let Some(transport) = self.unregister_transport(token) else { + let Some(transport) = self.deregister_transport(token) else { return Ok(()); }; - log::debug!(target: "reactor", token=token.0; "Handing over transport"); + log::debug!(target: "reactor", token=token.0; "Handing over transport {transport:?} with token {}", token.0); self.service.handover_transport(token, transport); } Action::Send(token, data) => { - log::trace!(target: "reactor", "Sending {} bytes to {token:?}", data.len()); + log::trace!(target: "reactor", token=token.0; "Sending {} bytes to {token:?}", data.len()); if let Some(transport) = self.transports.get_mut(&token) { if let Err(e) = transport.write_atomic(&data) { log::error!(target: "reactor", "Fatal error writing to transport {token:?}, disconnecting. Error details: {e:?}"); - if let Some(transport) = self.unregister_transport(token) { + if let Some(transport) = self.deregister_transport(token) { return Err(Error::TransportDisconnect(token, transport)); } } } else { - log::error!(target: "reactor", "Transport {token:?} is not in the reactor"); + log::error!(target: "reactor", token=token.0; "No transport with token {token:?} is known!"); } } Action::SetTimer(duration) => { @@ -562,27 +572,27 @@ impl Runtime { log::info!(target: "reactor", "Shutdown"); } - fn unregister_listener(&mut self, token: Token) -> Option { + fn deregister_listener(&mut self, token: Token) -> Option { let Some(mut source) = self.listeners.remove(&token) else { - log::warn!(target: "reactor", token=token.0; "Unregistering non-registered listener"); + log::warn!(target: "reactor", token=token.0; "Deregistering non-registered listener with token {}", token.0); return None; }; if let Err(err) = self.poll.registry().deregister(&mut source) { - log::warn!(target: "reactor", token=token.0; "Failed to deregister listener from mio: {err}"); + log::warn!(target: "reactor", token=token.0; "Failed to deregister listener with token {} from mio: {err}", token.0); } Some(source) } - fn unregister_transport(&mut self, token: Token) -> Option { + fn deregister_transport(&mut self, token: Token) -> Option { let Some(mut source) = self.transports.remove(&token) else { - log::warn!(target: "reactor", token=token.0; "Unregistering non-registered transport"); + log::warn!(target: "reactor", token=token.0; "Deregistering non-registered transport with token {}", token.0); return None; }; if let Err(err) = self.poll.registry().deregister(&mut source) { - log::warn!(target: "reactor", token=token.0; "Failed to deregister transport from mio: {err}"); + log::warn!(target: "reactor", token=token.0; "Failed to deregister transport with token {} from mio: {err}", token.0); } Some(source) diff --git a/crates/radicle-node/src/reactor/session.rs b/crates/radicle-node/src/reactor/session.rs index 38bcdcf1..74885e05 100644 --- a/crates/radicle-node/src/reactor/session.rs +++ b/crates/radicle-node/src/reactor/session.rs @@ -100,7 +100,7 @@ impl Display for ProtocolArtifact { } } -#[derive(Copy, Clone, Eq, PartialEq)] +#[derive(Copy, Clone, Eq, PartialEq, Debug)] pub struct Protocol { pub(crate) state: M, pub(crate) session: S, diff --git a/crates/radicle-node/src/runtime.rs b/crates/radicle-node/src/runtime.rs index 7b24c994..115ef587 100644 --- a/crates/radicle-node/src/runtime.rs +++ b/crates/radicle-node/src/runtime.rs @@ -1,6 +1,7 @@ pub mod handle; pub mod thread; +use std::fmt::Debug; use std::path::PathBuf; use std::{fs, io, net}; @@ -131,7 +132,11 @@ impl Runtime { signer: Device, ) -> Result where - G: crypto::signature::Signer + Ecdh + Clone + 'static, + G: crypto::signature::Signer + + Ecdh + + Clone + + Debug + + 'static, { let id = *signer.public_key(); let alias = config.alias.clone(); diff --git a/crates/radicle-node/src/test/node.rs b/crates/radicle-node/src/test/node.rs index f15109d7..2256b7cd 100644 --- a/crates/radicle-node/src/test/node.rs +++ b/crates/radicle-node/src/test/node.rs @@ -1,3 +1,4 @@ +use std::fmt::Debug; use std::io::BufRead as _; use std::mem::ManuallyDrop; use std::path::Path; @@ -455,7 +456,7 @@ impl Node { } } -impl + Signer + Clone> Node { +impl + Signer + Clone + Debug> Node { /// Spawn a node in its own thread. pub fn spawn(self) -> NodeHandle { let alias = self.config.alias.clone(); diff --git a/crates/radicle-node/src/wire.rs b/crates/radicle-node/src/wire.rs index 9aeded65..1cb4b330 100644 --- a/crates/radicle-node/src/wire.rs +++ b/crates/radicle-node/src/wire.rs @@ -3,6 +3,7 @@ //! We use the Noise XK handshake pattern to establish an encrypted stream with a remote peer. use std::collections::hash_map::Entry; use std::collections::VecDeque; +use std::fmt::Debug; use std::sync::Arc; use std::{io, net, time}; @@ -490,7 +491,7 @@ impl reactor::ReactionHandler for Wire where D: service::Store + Send, S: WriteStorage + Send + 'static, - G: crypto::signature::Signer + Ecdh + Clone + Send, + G: crypto::signature::Signer + Ecdh + Clone + Send + Debug, { type Listener = Listener; type Transport = Transport>;