use std::io::ErrorKind;
use std::net::TcpStream;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::Sender;
use std::sync::Arc;
use log::{debug, info};
use tungstenite::stream::MaybeTlsStream;
use tungstenite::{connect, Message, WebSocket};
use url::Url;
use crate::error::BinanceConnectError;
use crate::futures_usd::deserializer::deserialize;
use crate::futures_usd::enums::events::Event;
use crate::futures_usd::stream::WouldBlockConfig;
pub fn client(
sender: Sender<Event>,
url: Url,
stop_signal: Arc<AtomicBool>,
would_block_config: WouldBlockConfig,
subscribe_payload: Option<String>,
) -> Result<(), BinanceConnectError> {
let mut socket: WebSocket<MaybeTlsStream<TcpStream>> = socket(url)?;
if let Some(subscribe_payload) = subscribe_payload {
debug!("{:?}", subscribe_payload);
socket.send(Message::Text(subscribe_payload))?;
}
while !stop_signal.load(Ordering::Relaxed) {
match socket.read() {
Ok(message) => match message {
Message::Text(json_response) => {
if stop_signal.load(Ordering::Relaxed) {
return Ok(());
};
let event: Event = deserialize(json_response)?;
sender.send(event)?;
}
Message::Ping(ping) => {
if stop_signal.load(Ordering::Relaxed) {
return Ok(());
};
socket.send(Message::Pong(ping))?;
debug!("pong");
}
_ => {}
},
Err(err) => match err {
tungstenite::Error::Io(ref io_err) if io_err.kind() == ErrorKind::WouldBlock => {
if stop_signal.load(Ordering::Relaxed) {
return Ok(());
};
if would_block_config.error_on_block {
Err(BinanceConnectError::SocketError(err))?;
}
debug!(
"futures_usd client thread slept {:?} because of WouldBlock error",
would_block_config.time_out
);
std::thread::sleep(would_block_config.time_out);
}
_ => {
if stop_signal.load(Ordering::Relaxed) {
return Ok(());
};
Err(BinanceConnectError::SocketError(err))?;
}
},
}
}
Ok(())
}
fn socket(url: Url) -> Result<WebSocket<MaybeTlsStream<TcpStream>>, BinanceConnectError> {
let (socket, _) = connect(url)?;
Ok(socket)
}