use std::convert::From;
use std::borrow::Borrow;
use std::sync::mpsc;
use std::sync::mpsc::Sender;
use std::io::{Error, ErrorKind, Result};
use ws;
use ws::{CloseCode, Handshake, Handler, Message};
use thread::Joiner;
pub struct WebSocket {
ws_tx: ws::Sender,
_raii_joiner: Joiner,
}
impl WebSocket {
pub fn new<U: Borrow<str>>(url_borrow: U) -> Result<Self> {
let url = url_borrow.borrow().to_owned();
let (tx, rx) = mpsc::channel();
let joiner = thread!("WebSocketLogger", move || {
struct Client {
ws_tx: ws::Sender,
tx: Sender<Result<ws::Sender>>,
};
impl Handler for Client {
fn on_open(&mut self, _: Handshake) -> ws::Result<()> {
if self.tx.send(Ok(self.ws_tx.clone())).is_err() {
Err(ws::Error {
kind: ws::ErrorKind::Internal,
details: From::from("Channel error - Could not send ws_tx."),
})
} else {
Ok(())
}
}
}
let mut tx_opt = Some(tx.clone());
match ws::connect(url, |ws_tx| {
Client {
ws_tx: ws_tx,
tx: unwrap!(tx_opt.take(), "Logic Error! Report as bug."),
}
}) {
Ok(()) => (),
Err(e) => {
let _ = tx.send(Err(Error::new(ErrorKind::Other, format!("{:?}", e))));
}
}
});
match rx.recv() {
Ok(Ok(ws_tx)) => {
Ok(WebSocket {
ws_tx: ws_tx,
_raii_joiner: Joiner::new(joiner),
})
}
Ok(Err(e)) => Err(e),
Err(e) => Err(Error::new(ErrorKind::Other, format!("WebSocket Logger Error: {:?}", e))),
}
}
pub fn write_all(&self, buf: &[u8]) -> Result<()> {
self.ws_tx
.send(Message::Binary(buf.to_owned()))
.map_err(|e| Error::new(ErrorKind::Other, format!("{:?}", e)))
}
}
impl Drop for WebSocket {
fn drop(&mut self) {
let _ = self.ws_tx.close(CloseCode::Normal);
}
}