use {Event, RawFd, Result, Signal};
use util;
use std::ptr;
use std::slice;
use libc;
const MIN_BUFFER_SIZE: usize = 32;
const MAX_BUFFER_SIZE: usize = 1024 * 1024;
pub struct Mux {
kq: RawFd,
}
impl Mux {
pub fn new() -> Result<Self> {
let kq = unsafe { libc::kqueue() };
if kq < 0 {
return Err(util::errno());
}
Ok(Mux {
kq: kq as RawFd,
})
}
pub fn subscribe(&mut self, event: Event) -> Result<()> {
match event {
Event::Readable(fd) => {
push(self.kq, fd, libc::EVFILT_READ, libc::EV_ADD, 0)
},
Event::Writable(fd) => {
push(self.kq, fd, libc::EVFILT_WRITE, libc::EV_ADD, 0)
},
Event::Interrupt(signal) => {
let r = push(self.kq, signal, libc::EVFILT_SIGNAL, libc::EV_ADD, 0);
if r.is_ok() {
let _ = util::sigaction(signal, libc::SIG_IGN);
}
r
},
Event::Notified => {
push(self.kq, 0, libc::EVFILT_USER, libc::EV_ADD | libc::EV_CLEAR, 0)
},
}
}
pub fn unsubscribe(&mut self, event: Event) -> Result<()> {
match event {
Event::Readable(fd) => {
push(self.kq, fd, libc::EVFILT_READ, libc::EV_DELETE, 0)
},
Event::Writable(fd) => {
push(self.kq, fd, libc::EVFILT_WRITE, libc::EV_DELETE, 0)
},
Event::Interrupt(signal) => {
let r = push(self.kq, signal, libc::EVFILT_SIGNAL, libc::EV_DELETE, 0);
if r.is_ok() {
let _ = util::sigaction(signal, libc::SIG_DFL);
}
r
},
Event::Notified => {
push(self.kq, 0, libc::EVFILT_USER, libc::EV_DELETE, 0)
},
}
}
pub fn poll<'a>(&mut self, buf: &'a mut Buffer, timeout: Option<u64>) -> Result<Iter<'a>> {
let ts = if let Some(ns) = timeout {
&libc::timespec {
tv_sec: (ns / 1_000_000_000) as libc::time_t,
tv_nsec: (ns % 1_000_000_000) as libc::c_long,
}
} else {
ptr::null()
};
let n = unsafe {
let r = libc::kevent(self.kq, ptr::null(), 0, buf.slab.as_mut_ptr(),
buf.size as libc::c_int, ts);
if r < 0 {
return Err(util::errno());
}
r as usize
};
if n == buf.size && n < MAX_BUFFER_SIZE {
let _ = buf.slab.resize(n * 2);
}
Ok(Iter {
raw: unsafe { slice::from_raw_parts(buf.slab.as_ptr(), n) },
})
}
}
impl Drop for Mux {
fn drop(&mut self) {
unsafe {
libc::close(self.kq);
}
}
}
pub struct Buffer {
slab: util::Slab<libc::kevent>,
size: usize,
}
impl Buffer {
pub fn new() -> Result<Self> {
Ok(Buffer {
slab: try!(util::Slab::allocate(MIN_BUFFER_SIZE)),
size: MIN_BUFFER_SIZE,
})
}
}
pub struct Iter<'a> {
raw: &'a [libc::kevent],
}
impl<'a> Iterator for Iter<'a> {
type Item = Event;
fn next(&mut self) -> Option<Event> {
while self.raw.len() > 0 {
let raw = self.raw[0];
self.raw = &self.raw[1..];
match raw.filter {
libc::EVFILT_READ => return Some(Event::Readable(raw.ident as RawFd)),
libc::EVFILT_WRITE => return Some(Event::Writable(raw.ident as RawFd)),
libc::EVFILT_SIGNAL => return Some(Event::Interrupt(raw.ident as Signal)),
libc::EVFILT_USER => return Some(Event::Notified),
_ => {},
};
}
None
}
}
pub struct Notifier {
kq: RawFd,
}
impl Notifier {
pub fn attach(mux: &mut Mux) -> Result<Self> {
let kq = unsafe { libc::dup(mux.kq) };
if kq < 0 {
return Err(util::errno());
}
Ok(Notifier {
kq: kq,
})
}
pub fn notify(&mut self) -> Result<()> {
push(self.kq, 0, libc::EVFILT_USER, 0, libc::NOTE_TRIGGER)
}
}
impl Drop for Notifier {
fn drop(&mut self) {
unsafe {
libc::close(self.kq);
}
}
}
fn push(kq: RawFd, ident: libc::c_int, filter: libc::c_short,
flags: libc::c_ushort, fflags: libc::uint32_t) -> Result<()> {
let kevent = libc::kevent {
ident: ident as libc::uintptr_t,
filter: filter,
flags: flags,
fflags: fflags,
data: 0,
udata: ptr::null_mut(),
};
unsafe {
let r = libc::kevent(kq, &kevent, 1, ptr::null_mut(), 0, ptr::null());
if r < 0 {
Err(util::errno())
} else {
Ok(())
}
}
}