msgpacknet 0.1.0

A networking layer based on MessagePack messages.
Documentation
use std::sync::{Mutex, Arc, Condvar};
use std::collections::VecDeque;
#[cfg(feature = "nightly")]
use std::time::Duration;

pub struct QueueInner<T> {
    msgs: Mutex<(VecDeque<T>, bool)>,
    waiter: Condvar,
}

#[derive(Clone)]
pub struct Queue<T>(Arc<QueueInner<T>>);

impl<T> Queue<T> {
    pub fn new() -> Self {
        Queue(Arc::new(QueueInner{msgs: Mutex::new((VecDeque::new(), true)), waiter: Condvar::new()}))
    }

    pub fn get(&self) -> Option<T> {
        let mut lock = self.0.msgs.lock().expect("Lock poisoned");
        while lock.1 && lock.0.is_empty() {
            lock = self.0.waiter.wait(lock).expect("Lock poisoned");
        }
        if !lock.1 && lock.0.is_empty() {
            return None;
        }
        Some(lock.0.pop_front().unwrap())
    }

    #[cfg(feature = "nightly")]
    pub fn get_timeout(&self, timeout: Duration) -> Option<Option<T>> {
        let mut lock = self.0.msgs.lock().expect("Lock poisoned");
        while lock.1 && lock.0.is_empty() {
            let (new_lock, result) = self.0.waiter.wait_timeout(lock, timeout).expect("Lock poisoned");
            if result.timed_out() {
                return None;
            }
            lock = new_lock;
        }
        if !lock.1 && lock.0.is_empty() {
            return Some(None);
        }
        Some(Some(lock.0.pop_front().unwrap()))
    }

    pub fn close(&self) {
        self.0.msgs.lock().expect("Lock poisoned").1 = false;
    }

    pub fn put(&self, item: T) {
        let mut lock = self.0.msgs.lock().expect("Lock poisoned");
        if !lock.1 {
            return;
        }
        lock.0.push_back(item);
        self.0.waiter.notify_all();
    }
}

impl<T> Iterator for Queue<T> {
    type Item = T;

    fn next(&mut self) -> Option<Self::Item> {
        self.get()
    }
}