use super::{Queue, SyncQueue};
use time;
use std::{mem, ptr, ops, usize, u32};
use std::sync::{Arc, Mutex, MutexGuard, Condvar};
use std::sync::atomic::{self, AtomicUsize, Ordering};
pub struct LinkedQueue<T: Send> {
inner: Arc<QueueInner<T>>,
}
impl<T: Send> LinkedQueue<T> {
pub fn new() -> LinkedQueue<T> {
LinkedQueue::with_capacity(usize::MAX)
}
pub fn with_capacity(capacity: usize) -> LinkedQueue<T> {
LinkedQueue {
inner: Arc::new(QueueInner::new(capacity))
}
}
pub fn len(&self) -> usize {
self.inner.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn offer(&self, e: T) -> Result<(), T> {
self.inner.offer(e)
}
pub fn offer_ms(&self, e: T, ms: u32) -> Result<(), T> {
self.inner.offer_ms(e, ms)
}
pub fn put(&self, e: T) {
self.inner.put(e);
}
pub fn poll(&self) -> Option<T> {
self.inner.poll()
}
pub fn poll_ms(&self, ms: u32) -> Option<T> {
self.inner.poll_ms(ms)
}
pub fn take(&self) -> T {
self.inner.take()
}
}
impl<T: Send> Queue<T> for LinkedQueue<T> {
fn poll(&self) -> Option<T> {
LinkedQueue::poll(self)
}
fn is_empty(&self) -> bool {
LinkedQueue::is_empty(self)
}
fn offer(&self, e: T) -> Result<(), T> {
LinkedQueue::offer(self, e)
}
}
impl<T: Send> SyncQueue<T> for LinkedQueue<T> {
fn take(&self) -> T {
LinkedQueue::take(self)
}
fn put(&self, e: T) {
LinkedQueue::put(self, e)
}
}
impl<T: Send> Clone for LinkedQueue<T> {
fn clone(&self) -> LinkedQueue<T> {
LinkedQueue { inner: self.inner.clone() }
}
}
struct QueueInner<T: Send> {
capacity: usize,
count: AtomicUsize,
head: Mutex<NodePtr<T>>,
last: Mutex<NodePtr<T>>,
not_empty: Condvar,
not_full: Condvar,
}
impl<T: Send> QueueInner<T> {
fn new(capacity: usize) -> QueueInner<T> {
let head = NodePtr::new(Node::empty());
QueueInner {
capacity: capacity,
count: AtomicUsize::new(0),
head: Mutex::new(head),
last: Mutex::new(head),
not_empty: Condvar::new(),
not_full: Condvar::new(),
}
}
fn len(&self) -> usize {
self.count.load(Ordering::Relaxed)
}
fn put(&self, e: T) {
self.offer_ms(e, u32::MAX)
.ok().expect("something went wrong");
}
fn offer(&self, e: T) -> Result<(), T> {
if self.len() == self.capacity {
return Err(e);
}
self.offer_ms(e, 0)
}
fn offer_ms(&self, e: T, mut dur: u32) -> Result<(), T> {
let mut last = self.last.lock()
.ok().expect("something went wrong");
if self.len() == self.capacity {
let mut now = time::precise_time_ns();
loop {
if dur == 0 {
return Err(e);
}
last = self.not_full.wait_timeout_ms(last, dur)
.ok().expect("something went wrong").0;
if self.len() != self.capacity {
break;
}
let n = time::precise_time_ns();
let d = (n - now) / 1_000_000;
if d >= dur as u64 {
dur = 0;
} else {
dur -= d as u32;
now = n;
}
}
}
enqueue(Node::new(e), &mut last);
let cnt = self.count.fetch_add(1, Ordering::Release);
if cnt + 1 < self.capacity {
self.not_full.notify_one();
}
drop(last);
self.notify_not_empty();
Ok(())
}
fn take(&self) -> T {
self.poll_ms(u32::MAX)
.expect("something went wrong")
}
fn poll(&self) -> Option<T> {
if self.len() == 0 {
return None;
}
self.poll_ms(0)
}
fn poll_ms(&self, mut dur: u32) -> Option<T> {
let mut head = self.head.lock()
.ok().expect("something went wrong");
if self.len() == 0 {
let mut now = time::precise_time_ns();
loop {
if dur == 0 {
return None;
}
head = self.not_empty.wait_timeout_ms(head, dur)
.ok().expect("something went wrong").0;
if self.len() != 0 {
break;
}
let n = time::precise_time_ns();
let d = (n - now) / 1_000_000;
if d >= dur as u64 {
dur = 0;
} else {
dur -= d as u32;
now = n;
}
}
}
atomic::fence(Ordering::Acquire);
let val = dequeue(&mut head);
let cnt = self.count.fetch_sub(1, Ordering::Relaxed);
if cnt > 1 {
self.not_empty.notify_one();
}
drop(head);
if cnt == self.capacity {
self.notify_not_full();
}
Some(val)
}
fn notify_not_full(&self) {
let _l = self.last.lock()
.ok().expect("something went wrong");
self.not_full.notify_one();
}
fn notify_not_empty(&self) {
let _l = self.head.lock()
.ok().expect("something went wrong");
self.not_empty.notify_one();
}
}
impl<T: Send> Drop for QueueInner<T> {
fn drop(&mut self) {
while let Some(_) = self.poll() {
}
}
}
fn dequeue<T: Send>(mut head: &mut MutexGuard<NodePtr<T>>) -> T {
let h = **head;
let mut first = h.next;
**head = first;
h.free();
first.item.take().expect("item already consumed")
}
fn enqueue<T: Send>(node: Node<T>, mut last: &mut MutexGuard<NodePtr<T>>) {
let ptr = NodePtr::new(node);
last.next = ptr;
**last = ptr;
}
struct Node<T: Send> {
next: NodePtr<T>,
item: Option<T>,
}
impl<T: Send> Node<T> {
fn new(val: T) -> Node<T> {
Node {
next: NodePtr::null(),
item: Some(val),
}
}
fn empty() -> Node<T> {
Node {
next: NodePtr::null(),
item: None,
}
}
}
struct NodePtr<T: Send> {
ptr: *mut Node<T>,
}
impl<T: Send> NodePtr<T> {
fn new(node: Node<T>) -> NodePtr<T> {
NodePtr { ptr: unsafe { mem::transmute(Box::new(node)) }}
}
fn null() -> NodePtr<T> {
NodePtr { ptr: ptr::null_mut() }
}
fn free(self) {
let NodePtr { ptr } = self;
let _: Box<Node<T>> = unsafe { mem::transmute(ptr) };
}
}
impl<T: Send> ops::Deref for NodePtr<T> {
type Target = Node<T>;
fn deref(&self) -> &Node<T> {
unsafe { mem::transmute(self.ptr) }
}
}
impl<T: Send> ops::DerefMut for NodePtr<T> {
fn deref_mut(&mut self) -> &mut Node<T> {
unsafe { mem::transmute(self.ptr) }
}
}
impl<T: Send> Clone for NodePtr<T> {
fn clone(&self) -> NodePtr<T> {
NodePtr { ptr: self.ptr }
}
}
impl<T: Send> Copy for NodePtr<T> {}
unsafe impl<T: Send> Send for NodePtr<T> {}