use std::sync::{Arc,Mutex};
extern crate bus;
use bus::{Bus,BusReader};
#[cfg(test)]
mod tests;
pub trait TakesMessage<M> {
fn take_message(&mut self, t: &M);
}
pub struct TcWriter<M>
where M: Sync + Clone {
producer: Arc<Mutex<Bus<M>>>,
}
impl<M> TcWriter<M>
where M: Sync + Clone {
#[inline]
pub fn new(capacity: usize) -> Self {
TcWriter { producer: Arc::new(Mutex::new(Bus::new(capacity))), }
}
pub fn apply_change(&self, m: M) {
self.producer.lock().unwrap().broadcast(m)
}
pub fn try_apply_change(&self, m: M) -> Result<(), M> {
if let Err(m) = self.producer.lock().unwrap().try_broadcast(m) {
Err(m)
} else {
Ok(())
}
}
pub fn add_reader<T: TakesMessage<M>>(&self, init: T) -> TcReader<T, M> {
TcReader {
data: init,
consumer: self.producer.lock().unwrap().add_rx()
}
}
}
impl<M> Clone for TcWriter<M>
where M: Sync + Clone {
fn clone(&self) -> Self {
TcWriter { producer: self.producer.clone(), }
}
fn clone_from(&mut self, source: &Self) {
self.producer = source.producer.clone();
}
}
pub struct TcReader<T,M>
where T: TakesMessage<M> {
data: T,
consumer: BusReader<M>,
}
impl<T,M> TcReader<T,M>
where T: TakesMessage<M>,
M: Sync + Clone {
pub fn update(&mut self) -> usize {
let mut count = 0;
while let Ok(msg) = self.consumer.try_recv() {
self.apply_given(&msg);
count += 1
}
count
}
pub fn update_return(&mut self) -> Vec<M> {
let mut v = vec![];
while let Ok(msg) = self.consumer.try_recv() {
self.apply_given(&msg);
v.push(msg);
}
v
}
pub fn update_limited(&mut self, limit: usize) -> usize {
let mut count = 0;
for _ in 0..limit {
if let Ok(msg) = self.consumer.try_recv() {
self.apply_given(&msg);
count += 1;
} else { break }
}
count
}
pub fn update_return_limited(&mut self, limit: usize) -> Vec<M> {
let mut v = vec![];
for _ in 0..limit {
if let Ok(msg) = self.consumer.try_recv() {
self.apply_given(&msg);
v.push(msg);
} else { break }
}
v
}
pub fn into_inner(self) -> T {
self.data
}
#[inline]
fn apply_given(&mut self, msg: &M) {
self.data.take_message(&msg);
}
}
impl<T,M> std::ops::Deref for TcReader<T,M>
where T: TakesMessage<M> {
type Target = T;
fn deref(&self) -> &T {
&self.data
}
}
impl<T,M> std::ops::DerefMut for TcReader<T,M>
where T: TakesMessage<M> {
fn deref_mut(&mut self) -> &mut T {
&mut self.data
}
}