ev 0.1.0

Cross-platform event loop primitives for unix systems.
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 the `kevent` call used our whole buffer there were likely more
        // updates that we didn't have room for. Double the buffer's size every
        // time this happens.
        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);
        }
    }
}


/// Push a kqueue configuration change.
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(())
        }
    }
}