From 145164f7f466bb28e0a6dabeb7a44676f2c854ee Mon Sep 17 00:00:00 2001 From: hdu_willsky Date: Mon, 20 Jul 2026 01:50:08 +0800 Subject: [PATCH] fix: stop sync client terminal poll loops --- socketio/src/client/client.rs | 87 ++++++++++++++++++++++++++--------- 1 file changed, 65 insertions(+), 22 deletions(-) diff --git a/socketio/src/client/client.rs b/socketio/src/client/client.rs index fe924307..f5b0dcd9 100644 --- a/socketio/src/client/client.rs +++ b/socketio/src/client/client.rs @@ -213,35 +213,64 @@ impl Client { let mut self_clone = self.clone(); // Use thread to consume items in iterator in order to call callbacks std::thread::spawn(move || { - // tries to restart a poll cycle whenever a 'normal' error occurs, - // it just panics on network errors, in case the poll cycle returned - // `Result::Ok`, the server receives a close frame so it's safe to - // terminate + // Reconnect only when explicitly enabled. Otherwise any terminal poll state must end + // this worker so a failed connection cannot leave a busy loop behind. for packet in self_clone.iter() { - let should_reconnect = match packet { - Err(Error::IncompleteResponseFromEngineIo(_)) => { - //TODO: 0.3.X handle errors - //TODO: logging error - true - } - Ok(Packet { + let action = match &packet { + Err(_) + | Ok(Packet { packet_type: PacketId::Disconnect, .. - }) => match self_clone.builder.lock() { - Ok(builder) => builder.reconnect_on_disconnect, - Err(_) => false, - }, - _ => false, + }) => { + let (reconnect_enabled, reconnect_on_disconnect) = + match self_clone.builder.lock() { + Ok(builder) => (builder.reconnect, builder.reconnect_on_disconnect), + Err(_) => (false, false), + }; + poll_action(&packet, reconnect_enabled, reconnect_on_disconnect) + } + _ => PollAction::Continue, }; - if should_reconnect { - let _ = self_clone.disconnect(); - let _ = self_clone.reconnect(); + match action { + PollAction::Continue => {} + PollAction::Reconnect => { + let _ = self_clone.disconnect(); + let _ = self_clone.reconnect(); + } + PollAction::Stop => break, } } }); } } +#[derive(Debug, PartialEq, Eq)] +enum PollAction { + Continue, + Reconnect, + Stop, +} + +fn poll_action( + packet: &Result, + reconnect_enabled: bool, + reconnect_on_disconnect: bool, +) -> PollAction { + match packet { + Err(Error::IncompleteResponseFromEngineIo(_)) if reconnect_enabled => PollAction::Reconnect, + Err(_) => PollAction::Stop, + Ok(Packet { + packet_type: PacketId::Disconnect, + .. + }) if reconnect_enabled && reconnect_on_disconnect => PollAction::Reconnect, + Ok(Packet { + packet_type: PacketId::Disconnect, + .. + }) => PollAction::Stop, + _ => PollAction::Continue, + } +} + pub(crate) struct Iter { socket: Arc>, } @@ -255,9 +284,9 @@ impl Iterator for Iter { Ok(socket) => match socket.poll() { Err(err) => Some(Err(err)), Ok(Some(packet)) => Some(Ok(packet)), - // If the underlying engineIO connection is closed, - // throw an error so we know to reconnect - Ok(None) => Some(Err(Error::StoppedEngineIoSocket)), + // A stopped Engine.IO socket is terminal. Returning another item here makes the + // callback worker poll the already-closed transport forever. + Ok(None) => None, }, Err(_) => { // Lock is poisoned, our iterator is useless. @@ -282,6 +311,20 @@ mod test { use std::time::{Duration, SystemTime}; use url::Url; + #[test] + fn terminal_polling_states_stop_without_reconnect() { + let stopped = Err(Error::StoppedEngineIoSocket); + let transport_error = Err(Error::IncompleteResponseFromEngineIo( + rust_engineio::Error::IncompletePacket(), + )); + + assert_eq!(poll_action(&stopped, false, false), PollAction::Stop); + assert_eq!( + poll_action(&transport_error, false, false), + PollAction::Stop + ); + } + #[test] #[serial(reconnect)] fn socket_io_reconnect_integration() -> Result<()> {