use lossyq::spsc::*;
use scheduler;
use super::{wrap};
use super::observer::{CountingReporter, TaskTracer};
use super::super::{Message, Schedule, TaskState, Error};
use super::super::elem::{source};
use std::sync::atomic::{AtomicUsize, Ordering};
struct TestSource {
ret: Schedule,
exec_count: usize,
}
impl source::Source for TestSource {
type OutputType = u64;
fn process(
&mut self,
_output: &mut Sender<Message<Self::OutputType>>)
-> Schedule {
self.exec_count += 1;
println!("exec count: {}",self.exec_count);
self.ret
}
}
#[test]
fn sched_add_task() {
let mut sched = scheduler::new();
let first_id : usize;
{
let (source_task, mut _source_out) =
source::new( "Source", 2, Box::new(TestSource{ret:Schedule::DelayUSec(2_000), exec_count:0}));
let result = sched.add_task(source_task);
assert!(result.is_ok());
first_id = match result {
Ok(task_id) => { task_id.id() }
_ => { 9999 }
};
assert!(first_id != 9999);
}
{
let (source_task, mut _source_out) =
source::new( "Source", 2, Box::new(TestSource{ret:Schedule::DelayUSec(2_000), exec_count:0}));
let result = sched.add_task(source_task);
assert!(result.is_err());
let already_exists = match result {
Err(Error::AlreadyExists) => { true },
_ => { false }
};
assert!(already_exists);
}
{
let (source_task, mut _source_out) =
source::new( "Source 3", 2, Box::new(TestSource{ret:Schedule::DelayUSec(2_000), exec_count:0}));
let result = sched.add_task(source_task);
assert!(result.is_ok());
let third_id = match result {
Ok(task_id) => { task_id.id() }
_ => { 9999 }
};
assert!(third_id != first_id);
assert!(third_id != 9999);
}
}
#[test]
fn wrap_execute_time_delayed() {
let (source_task, mut _source_out) =
source::new( "Source", 2, Box::new(TestSource{ret:Schedule::DelayUSec(2_000), exec_count:0}));
let mut wrp = wrap::new(source_task, 99);
let mut obs = CountingReporter::new();
let tim = AtomicUsize::new(0);
assert_eq!(wrp.execute(&mut obs, &tim), TaskState::TimeWait(2_000));
assert_eq!(obs.executed, 1);
assert_eq!(obs.delayed, 0);
assert_eq!(obs.time_wait, 0);
assert_eq!(wrp.execute(&mut obs, &tim), TaskState::TimeWait(2_000));
assert_eq!(obs.executed, 1);
assert_eq!(obs.delayed, 1);
assert_eq!(obs.time_wait, 1);
tim.fetch_add(2_001, Ordering::SeqCst);
assert_eq!(wrp.execute(&mut obs, &tim), TaskState::TimeWait(4_001));
assert_eq!(obs.executed, 2);
assert_eq!(obs.delayed, 1);
assert_eq!(obs.time_wait, 1);
assert_eq!(wrp.execute(&mut obs, &tim), TaskState::TimeWait(4_001));
assert_eq!(obs.executed, 2);
assert_eq!(obs.delayed, 2);
assert_eq!(obs.time_wait, 2);
}
#[test]
fn wrap_execute_traced() {
let (source_task, mut _source_out) =
source::new( "Source", 2, Box::new(TestSource{ret:Schedule::DelayUSec(2_000), exec_count:0}));
let mut wrp = wrap::new(source_task, 99);
let mut obs = TaskTracer::new();
let tim = AtomicUsize::new(0);
wrp.execute(&mut obs, &tim);
wrp.execute(&mut obs, &tim);
wrp.execute(&mut obs, &tim);
tim.fetch_add(2_001, Ordering::SeqCst);
wrp.execute(&mut obs, &tim);
wrp.execute(&mut obs, &tim);
wrp.execute(&mut obs, &tim);
assert_eq!(true, false);
}