use crate::{
Interval, Job, SyncJob,
timeprovider::{DefaultTimeProvider, TimeProvider, local_offset}
};
use std::{
default::Default,
sync::{
Arc,
atomic::{AtomicBool, Ordering}
},
thread,
time::Duration
};
use time::UtcOffset;
#[derive(Debug)]
pub struct Scheduler<TP = DefaultTimeProvider> {
jobs: Vec<SyncJob<TP>>,
utc_offset: UtcOffset
}
impl Default for Scheduler {
fn default() -> Self {
Self::new()
}
}
impl Scheduler {
pub fn new() -> Self {
Self::new_with_utc_offset(local_offset())
}
pub const fn new_with_utc_offset(utc_offset: UtcOffset) -> Self {
Self::new_with_provider(utc_offset)
}
}
impl<TP: TimeProvider> Scheduler<TP> {
pub const fn new_with_provider(utc_offset: UtcOffset) -> Self {
Self {
jobs: Vec::new(),
utc_offset
}
}
pub fn every(&mut self, ival: Interval) -> &mut SyncJob<TP> {
let job = SyncJob::new(ival, self.utc_offset);
self.jobs.push(job);
let last_index = self.jobs.len() - 1;
&mut self.jobs[last_index]
}
pub fn run_pending(&mut self) {
let now = TP::now(self.utc_offset);
for job in &mut self.jobs {
if job.is_pending(now) {
job.execute(now);
}
}
}
}
impl Scheduler {
#[must_use = "The scheduler is halted when the returned handle is dropped"]
pub fn watch_thread(self, frequency: Duration) -> ScheduleHandle {
let stop = Arc::new(AtomicBool::new(false));
let my_stop = stop.clone();
let mut me = self;
let handle = thread::spawn(move || {
while !stop.load(Ordering::SeqCst) {
me.run_pending();
thread::sleep(frequency);
}
});
ScheduleHandle {
stop: my_stop,
thread_handle: Some(handle)
}
}
}
pub struct ScheduleHandle {
stop: Arc<AtomicBool>,
thread_handle: Option<thread::JoinHandle<()>>
}
impl ScheduleHandle {
pub fn stop(self) {}
}
impl Drop for ScheduleHandle {
fn drop(&mut self) {
self.stop.store(true, Ordering::SeqCst);
let handle = self.thread_handle.take();
handle.unwrap().join().ok();
}
}
#[cfg(test)]
mod tests {
use super::{Job, Scheduler, TimeProvider};
use crate::intervals::*;
use std::sync::{
Arc,
atomic::{AtomicU32, Ordering}
};
use time::{OffsetDateTime, UtcOffset, format_description::well_known::Rfc3339};
macro_rules! make_time_provider {
($name:ident : $($time:literal),+) => {
#[derive(Debug)]
struct $name;
static TIMES_TIME_REQUESTED: AtomicU32 = AtomicU32::new(0);
impl TimeProvider for $name {
fn now(utc_offset: UtcOffset) -> OffsetDateTime {
static TIMES: &[&str] = &[$($time),+];
let idx = TIMES_TIME_REQUESTED.fetch_add(1, Ordering::SeqCst) as usize;
OffsetDateTime::parse(TIMES[idx], &Rfc3339).unwrap().to_offset(utc_offset)
}
}
};
}
#[test]
fn test_every_plus() {
make_time_provider!(FakeTimeProvider :
"2019-10-22T12:40:00Z",
"2019-10-22T12:40:00Z",
"2019-10-22T12:50:20Z",
"2019-10-22T12:50:30Z"
);
let mut scheduler =
Scheduler::<FakeTimeProvider>::new_with_provider(UtcOffset::UTC);
let times_called = Arc::new(AtomicU32::new(0));
{
let times_called = times_called.clone();
scheduler
.every(10.minutes())
.plus(5.seconds())
.run(move || {
times_called.fetch_add(1, Ordering::SeqCst);
});
}
assert_eq!(1, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(0, times_called.load(Ordering::SeqCst));
assert_eq!(2, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(1, times_called.load(Ordering::SeqCst));
assert_eq!(3, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(1, times_called.load(Ordering::SeqCst));
assert_eq!(4, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
}
#[test]
fn test_every_at() {
make_time_provider!(FakeTimeProvider:
"2019-10-22T12:40:00Z",
"2019-10-22T12:40:10Z",
"2019-10-25T12:50:20Z",
"2019-10-25T15:23:30Z",
"2019-10-26T15:50:30Z"
);
let mut scheduler =
Scheduler::<FakeTimeProvider>::new_with_provider(UtcOffset::UTC);
let times_called = Arc::new(AtomicU32::new(0));
{
let times_called = times_called.clone();
scheduler.every(3.days()).at("15:23").run(move || {
times_called.fetch_add(1, Ordering::SeqCst);
});
}
assert_eq!(1, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(0, times_called.load(Ordering::SeqCst));
assert_eq!(2, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(1, times_called.load(Ordering::SeqCst));
assert_eq!(3, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(1, times_called.load(Ordering::SeqCst));
assert_eq!(4, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
}
#[test]
fn test_every_and_every() {
make_time_provider!(FakeTimeProvider:
"2019-10-22T12:40:01Z",
"2019-10-22T12:40:01Z",
"2019-10-22T12:40:02Z",
"2019-10-22T12:40:03Z",
"2019-10-22T12:40:04Z",
"2019-10-22T12:40:05Z",
"2019-10-22T12:40:06Z"
);
let mut scheduler =
Scheduler::<FakeTimeProvider>::new_with_provider(UtcOffset::UTC);
let times_called = Arc::new(AtomicU32::new(0));
{
let times_called = times_called.clone();
scheduler
.every(5.seconds())
.and_every(2.seconds())
.run(move || {
times_called.fetch_add(1, Ordering::SeqCst);
});
}
assert_eq!(1, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(2, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(0, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(3, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(1, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(4, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(1, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(5, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(2, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(6, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(3, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(7, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(4, times_called.load(Ordering::SeqCst));
}
#[test]
fn test_once() {
make_time_provider!(FakeTimeProvider:
"2019-10-22T12:40:01Z",
"2019-10-22T12:40:01Z",
"2019-10-22T12:40:02Z",
"2019-10-22T12:40:03Z"
);
let mut scheduler =
Scheduler::<FakeTimeProvider>::new_with_provider(UtcOffset::UTC);
let times_called = Arc::new(AtomicU32::new(0));
{
let times_called = times_called.clone();
scheduler.every(1.seconds()).once().run(move || {
times_called.fetch_add(1, Ordering::SeqCst);
});
}
assert_eq!(1, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(2, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(0, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(3, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(1, times_called.load(Ordering::SeqCst));
scheduler.run_pending();
assert_eq!(4, TIMES_TIME_REQUESTED.load(Ordering::SeqCst));
assert_eq!(1, times_called.load(Ordering::SeqCst));
}
}