extern crate rand;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::thread;
use crossbeam_queue::SegQueue;
use rand::Rng;
use crate::errors::Result;
use crate::Record;
pub struct Recordset {
instances: AtomicUsize,
record_queue_count: AtomicUsize,
record_queue_size: AtomicUsize,
record_queue: SegQueue<Result<Record>>,
active: AtomicBool,
task_id: AtomicUsize,
}
impl Recordset {
#[doc(hidden)]
pub fn new(rec_queue_size: usize, nodes: usize) -> Self {
let mut rng = rand::thread_rng();
let task_id = rng.gen::<usize>();
Recordset {
instances: AtomicUsize::new(nodes),
record_queue_size: AtomicUsize::new(rec_queue_size),
record_queue_count: AtomicUsize::new(0),
record_queue: SegQueue::new(),
active: AtomicBool::new(true),
task_id: AtomicUsize::new(task_id),
}
}
pub fn close(&self) {
self.active.store(false, Ordering::Relaxed)
}
pub fn is_active(&self) -> bool {
self.active.load(Ordering::Relaxed)
}
#[doc(hidden)]
pub fn push(&self, record: Result<Record>) -> Option<Result<Record>> {
if self.record_queue_count.fetch_add(1, Ordering::Relaxed)
< self.record_queue_size.load(Ordering::Relaxed)
{
self.record_queue.push(record);
return None;
}
self.record_queue_count.fetch_sub(1, Ordering::Relaxed);
Some(record)
}
pub fn task_id(&self) -> u64 {
self.task_id.load(Ordering::Relaxed) as u64
}
#[doc(hidden)]
pub fn signal_end(&self) {
if self.instances.fetch_sub(1, Ordering::Relaxed) == 1 {
self.close()
};
}
}
impl<'a> Iterator for &'a Recordset {
type Item = Result<Record>;
fn next(&mut self) -> Option<Result<Record>> {
loop {
if self.is_active() || !self.record_queue.is_empty() {
let result = self.record_queue.pop().ok();
if result.is_some() {
self.record_queue_count.fetch_sub(1, Ordering::Relaxed);
return result;
}
thread::yield_now();
continue;
} else {
return None;
}
}
}
}