use std::sync;
use std::io;
use std::collections;
use futures;
struct SemaphoreInner {
capacity: usize,
waiters: collections::VecDeque<futures::task::Task>,
}
pub struct Semaphore {
inner: sync::RwLock<SemaphoreInner>,
}
impl Semaphore {
pub fn new(initial: usize) -> Semaphore {
Semaphore {
inner: sync::RwLock::new(SemaphoreInner {
capacity: initial,
waiters: collections::VecDeque::new(),
}),
}
}
pub fn acquire(&self) -> SemaphoreHandle {
let mut lock_result = self.inner.write();
match lock_result {
Ok(ref mut guard) => {
if guard.capacity > 0 {
guard.capacity -= 1;
SemaphoreHandle::Completed(futures::future::result(Ok(())))
} else {
guard.waiters.push_back(futures::task::current());
SemaphoreHandle::Waiting
}
}
Err(err) => panic!("Lock failure {:?}", err),
}
}
pub fn release(&self) {
let mut lock_result = self.inner.write();
match lock_result {
Ok(ref mut guard) => {
if !guard.waiters.is_empty() {
guard.waiters.pop_front().unwrap().notify();
} else {
guard.capacity += 1;
}
}
Err(err) => panic!("Lock failure {:?}", err),
}
}
pub fn current_capacity(&self) -> usize {
let mut lock_result = self.inner.read();
match lock_result {
Ok(ref mut guard) => guard.capacity,
Err(_) => panic!("Lock failure"),
}
}
}
pub enum SemaphoreHandle {
Waiting,
Completed(futures::future::FutureResult<(), io::Error>),
}
impl futures::Future for SemaphoreHandle {
type Item = ();
type Error = io::Error;
fn poll(&mut self) -> Result<futures::Async<()>, io::Error> {
match self {
&mut SemaphoreHandle::Completed(_) => Ok(futures::Async::Ready(())),
&mut SemaphoreHandle::Waiting => {
*self = SemaphoreHandle::Completed(futures::future::result(Ok(())));
Ok(futures::Async::NotReady)
}
}
}
}