use futures::sync::mpsc as futures_mpsc;
use futures::{Async,Sink,Stream};
use futures::sink::Wait;
use std::io;
use std::os::raw::{c_int};
use std::sync::mpsc as std_mpsc;
use std::thread;
use std::time::Duration;
use tokio_core::reactor::{Handle,Remote};
use std::cell::UnsafeCell;
use remote::GetRemote;
#[derive(Clone,Copy,PartialEq,Eq,Debug)]
enum PollRequest {
Poll,
Close,
}
struct SelectFdRead {
fd: c_int,
read_fds: libc::fd_set,
}
impl SelectFdRead {
pub fn new(fd: c_int) -> Self {
use std::mem::uninitialized;
let mut read_fds : libc::fd_set = unsafe { uninitialized() };
unsafe { libc::FD_ZERO(&mut read_fds) };
SelectFdRead{
fd: fd,
read_fds: read_fds,
}
}
pub fn select(&mut self, timeout: Option<Duration>) -> bool {
use std::ptr::null_mut;
let mut timeout = timeout.map(|timeout|
libc::timeval{
tv_sec: timeout.as_secs() as libc::c_long,
tv_usec: (timeout.subsec_nanos() / 1000) as libc::c_long,
}
);
unsafe {
libc::FD_SET(self.fd, &mut self.read_fds);
libc::select(
self.fd+1,
&mut self.read_fds,
null_mut(),
null_mut(),
timeout.as_mut().map(|x| x as *mut _).unwrap_or(null_mut()),
);
libc::FD_ISSET(self.fd, &mut self.read_fds)
}
}
}
struct Inner {
fd: c_int,
_thread: thread::JoinHandle<()>,
pending_request: bool,
send_request: std_mpsc::SyncSender<PollRequest>,
send_response: Wait<futures_mpsc::Sender<()>>,
recv_response: futures_mpsc::Receiver<()>,
remote: Remote,
}
impl Inner {
fn poll_read(&mut self) -> Async<()> {
debug!("poll read");
if !self.pending_request {
let mut read_fds = SelectFdRead::new(self.fd);
if read_fds.select(Some(Duration::from_millis(0))) {
debug!("poll read: local ready");
return Async::Ready(());
} else {
debug!("poll read: not ready, start thread");
self.send_request.send(PollRequest::Poll).expect("select thread terminated");
self.pending_request = true;
}
}
match self.recv_response.poll().unwrap() {
Async::Ready(None) => unreachable!(),
Async::Ready(Some(())) => {
debug!("poll read: thread ready");
self.pending_request = false;
Async::Ready(())
},
Async::NotReady => {
debug!("poll read: thread not ready");
Async::NotReady
},
}
}
fn need_read(&mut self) {
match self.recv_response.poll().unwrap() {
Async::Ready(None) => unreachable!(),
Async::Ready(Some(())) => {
assert!(self.pending_request);
match self.recv_response.poll().unwrap() {
Async::Ready(None) => unreachable!(),
Async::Ready(Some(())) => unreachable!(),
Async::NotReady => (),
}
self.send_response.send(()).unwrap();
},
Async::NotReady => {
let mut read_fds = SelectFdRead::new(self.fd);
self.pending_request = true;
if read_fds.select(Some(Duration::from_millis(0))) {
self.send_response.send(()).unwrap();
} else {
debug!("poll need read: not ready, start thread");
self.send_request.send(PollRequest::Poll).expect("select thread terminated");
}
},
}
}
}
pub struct PollReadFd(UnsafeCell<Inner>);
impl PollReadFd {
pub fn new(fd: c_int, handle: &Handle) -> io::Result<Self> {
let (send_request, recv_request) = std_mpsc::sync_channel(1);
let (send_response, recv_response) = futures_mpsc::channel(1);
let outer_send_response = send_response.clone().wait();
let thread = thread::spawn(move || {
let mut read_fds = SelectFdRead::new(fd);
let mut send_response = send_response.wait();
loop {
debug!("[select thread] waiting for request");
match recv_request.recv() {
Ok(PollRequest::Poll) => (),
Ok(PollRequest::Close) => return,
Err(_) => return,
}
debug!("[select thread] start polling");
while !read_fds.select(Some(Duration::from_millis(1000))) {
match recv_request.try_recv() {
Ok(PollRequest::Poll) => unreachable!(),
Ok(PollRequest::Close) => return,
Err(_) => (), }
}
debug!("[select thread] read event");
if send_response.send(()).is_err() { return; }
}
});
Ok(PollReadFd(UnsafeCell::new(Inner{
fd: fd,
_thread: thread,
pending_request: false,
send_request: send_request,
send_response: outer_send_response,
recv_response: recv_response,
remote: handle.remote().clone(),
})))
}
fn inner(&self) -> &mut Inner {
unsafe { &mut *self.0.get() }
}
pub fn poll_read(&self) -> Async<()> {
self.inner().poll_read()
}
pub fn need_read(&self) {
self.inner().need_read()
}
}
impl GetRemote for PollReadFd {
fn remote(&self) -> &Remote {
&self.inner().remote
}
}
impl Drop for PollReadFd {
fn drop(&mut self) {
let _ = self.inner().send_request.send(PollRequest::Close);
}
}
#[cfg(windows)]
mod libc {
pub use libc::{c_int,c_uint,c_long};
pub use winapi::{timeval,fd_set,FD_SETSIZE,SOCKET};
pub use ws2_32::select;
pub unsafe fn FD_ZERO(set: *mut fd_set) {
let set = &mut *set;
set.fd_count = 0;
}
pub unsafe fn FD_SET(fd: c_int, set: *mut fd_set) {
if FD_ISSET(fd, set) { return; }
let set = &mut *set;
let fd = fd as c_uint as SOCKET;
if (set.fd_count as usize) < FD_SETSIZE {
set.fd_array[set.fd_count as usize] = fd;
set.fd_count += 1;
}
}
pub unsafe fn FD_ISSET(fd: c_int, set: *mut fd_set) -> bool {
let set = &mut *set;
let fd = fd as c_uint as SOCKET;
set.fd_array[..set.fd_count as usize].iter().any(|i| *i == fd)
}
}
#[cfg(unix)]
use libc;