use std::sync::mpsc;
use std::io;
use std::time::Duration;
use super::*;
use reactor;
use core::{SocketId, Message, PollReq};
use core::socket::{Request, Reply};
use core::config::ConfigOption;
use core;
use io_error::*;
#[doc(hidden)]
pub type ReplyReceiver = mpsc::Receiver<Reply>;
#[doc(hidden)]
pub struct RequestSender {
req_tx: EventLoopRequestSender,
socket_id: SocketId
}
impl RequestSender {
pub fn new(tx: EventLoopRequestSender, id: SocketId) -> RequestSender {
RequestSender {
req_tx: tx,
socket_id: id
}
}
fn child_sender(&self, eid: core::EndpointId) -> endpoint::RequestSender {
endpoint::RequestSender::new(self.req_tx.clone(), self.socket_id, eid)
}
fn send(&self, req: Request) -> io::Result<()> {
self.req_tx.send(reactor::Request::Socket(self.socket_id, req)).map_err(from_send_error)
}
}
pub struct Socket {
request_sender: RequestSender,
reply_receiver: ReplyReceiver
}
impl Socket {
#[doc(hidden)]
pub fn new(request_tx: RequestSender, reply_rx: ReplyReceiver) -> Socket {
Socket {
request_sender: request_tx,
reply_receiver: reply_rx
}
}
#[doc(hidden)]
pub fn id(&self) -> SocketId {
self.request_sender.socket_id
}
pub fn create_poll_req(&self, recv: bool, send: bool) -> PollReq {
PollReq {
sid: self.id(),
recv: recv,
send: send
}
}
pub fn connect(&mut self, url: &str) -> io::Result<endpoint::Endpoint> {
let request = Request::Connect(From::from(url));
self.call(request, |reply| self.on_connect_reply(reply))
}
fn on_connect_reply(&self, reply: Reply) -> io::Result<endpoint::Endpoint> {
match reply {
Reply::Connect(id) => {
let request_tx = self.request_sender.child_sender(id);
let ep = endpoint::Endpoint::new(request_tx, true);
Ok(ep)
},
Reply::Err(e) => Err(e),
_ => self.unexpected_reply()
}
}
pub fn bind(&mut self, url: &str) -> io::Result<endpoint::Endpoint> {
let request = Request::Bind(From::from(url));
self.call(request, |reply| self.on_bind_reply(reply))
}
fn on_bind_reply(&self, reply: Reply) -> io::Result<endpoint::Endpoint> {
match reply {
Reply::Bind(id) => {
let request_tx = self.request_sender.child_sender(id);
let ep = endpoint::Endpoint::new(request_tx, false);
Ok(ep)
},
Reply::Err(e) => Err(e),
_ => self.unexpected_reply()
}
}
pub fn send(&mut self, buffer: Vec<u8>) -> io::Result<()> {
self.send_msg(Message::from_body(buffer))
}
pub fn send_msg(&mut self, msg: Message) -> io::Result<()> {
let request = Request::Send(msg, false);
self.call(request, |reply| self.on_send_reply(reply))
}
pub fn try_send(&mut self, buffer: Vec<u8>) -> io::Result<()> {
self.try_send_msg(Message::from_body(buffer))
}
pub fn try_send_msg(&mut self, msg: Message) -> io::Result<()> {
let request = Request::Send(msg, true);
self.call(request, |reply| self.on_send_reply(reply))
}
fn on_send_reply(&self, reply: Reply) -> io::Result<()> {
match reply {
Reply::Send => Ok(()),
Reply::Err(e) => Err(e),
_ => self.unexpected_reply()
}
}
pub fn recv(&mut self) -> io::Result<Vec<u8>> {
self.recv_msg().map(|msg| msg.into())
}
pub fn recv_msg(&mut self) -> io::Result<Message> {
let request = Request::Recv(false);
self.call(request, |reply| self.on_recv_reply(reply))
}
pub fn try_recv(&mut self) -> io::Result<Vec<u8>> {
self.try_recv_msg().map(|msg| msg.into())
}
pub fn try_recv_msg(&mut self) -> io::Result<Message> {
let request = Request::Recv(true);
self.call(request, |reply| self.on_recv_reply(reply))
}
fn on_recv_reply(&self, reply: Reply) -> io::Result<Message> {
match reply {
Reply::Recv(msg) => Ok(msg),
Reply::Err(e) => Err(e),
_ => self.unexpected_reply()
}
}
pub fn set_send_timeout(&mut self, timeout: Option<Duration>) -> io::Result<()> {
self.set_option(ConfigOption::SendTimeout(timeout))
}
pub fn set_recv_timeout(&mut self, timeout: Option<Duration>) -> io::Result<()> {
self.set_option(ConfigOption::RecvTimeout(timeout))
}
pub fn set_send_priority(&mut self, priority: u8) -> io::Result<()> {
self.set_option(ConfigOption::SendPriority(priority))
}
pub fn set_recv_priority(&mut self, priority: u8) -> io::Result<()> {
self.set_option(ConfigOption::RecvPriority(priority))
}
pub fn set_tcp_nodelay(&mut self, value: bool) -> io::Result<()> {
self.set_option(ConfigOption::TcpNoDelay(value))
}
pub fn set_option(&mut self, cfg_opt: ConfigOption) -> io::Result<()> {
let request = Request::SetOption(cfg_opt);
self.call(request, |reply| self.on_set_option_reply(reply))
}
fn on_set_option_reply(&self, reply: Reply) -> io::Result<()> {
match reply {
Reply::SetOption => Ok(()),
Reply::Err(e) => Err(e),
_ => self.unexpected_reply()
}
}
fn call<T, F : FnOnce(Reply) -> io::Result<T>>(&self, request: Request, process: F) -> io::Result<T> {
self.execute_request(request).and_then(process)
}
fn execute_request(&self, request: Request) -> io::Result<Reply> {
self.send_request(request).and_then(|_| self.recv_reply())
}
fn send_request(&self, request: Request) -> io::Result<()> {
self.request_sender.send(request)
}
fn recv_reply(&self) -> io::Result<Reply> {
self.reply_receiver.receive()
}
fn unexpected_reply<T>(&self) -> io::Result<T> {
Err(other_io_error("unexpected reply"))
}
}
impl Drop for Socket {
fn drop(&mut self) {
let _ = self.send_request(Request::Close);
let _ = self.recv_reply();
}
}