Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 65 additions & 22 deletions socketio/src/client/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Packet>,
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<RwLock<RawClient>>,
}
Expand All @@ -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.
Expand All @@ -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<()> {
Expand Down