use crate::schedule::compiled::CompiledSchedule;
use crate::schedule::spec::{OverlapPolicy, ScheduleOnFailure};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TickAction {
Dispatch,
Skip,
Queue,
ForbidAbort,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RunOutcome {
Success,
Failure,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AfterRun {
Continue { dispatch_pending: bool },
ExitOk,
ExitFailure { consecutive: u64 },
}
pub struct SchedulerState {
overlap: OverlapPolicy,
on_failure: ScheduleOnFailure,
max_runs: Option<u64>,
max_consecutive_failures: Option<u64>,
successful_runs: u64,
consecutive_failures: u64,
pending: bool,
}
impl SchedulerState {
pub fn new(c: &CompiledSchedule) -> Self {
Self {
overlap: c.overlap_policy,
on_failure: c.on_failure,
max_runs: c.max_runs,
max_consecutive_failures: c.max_consecutive_failures,
successful_runs: 0,
consecutive_failures: 0,
pending: false,
}
}
pub fn consecutive_failures(&self) -> u64 {
self.consecutive_failures
}
pub fn on_tick(&mut self, running: bool) -> TickAction {
if !running {
return TickAction::Dispatch;
}
match self.overlap {
OverlapPolicy::Skip => TickAction::Skip,
OverlapPolicy::Queue => {
self.pending = true;
TickAction::Queue
}
OverlapPolicy::Forbid => TickAction::ForbidAbort,
}
}
pub fn on_run_finished(&mut self, outcome: RunOutcome) -> AfterRun {
match outcome {
RunOutcome::Success => {
self.consecutive_failures = 0;
self.successful_runs += 1;
if let Some(max) = self.max_runs
&& self.successful_runs >= max
{
return AfterRun::ExitOk;
}
}
RunOutcome::Failure => {
self.consecutive_failures += 1;
if matches!(self.on_failure, ScheduleOnFailure::Stop) {
return AfterRun::ExitFailure {
consecutive: self.consecutive_failures,
};
}
if let Some(max) = self.max_consecutive_failures
&& self.consecutive_failures >= max
{
return AfterRun::ExitFailure {
consecutive: self.consecutive_failures,
};
}
}
}
let dispatch_pending = self.pending;
self.pending = false;
AfterRun::Continue { dispatch_pending }
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::schedule::spec::ScheduleSpec;
fn state(yaml: &str) -> SchedulerState {
let spec: ScheduleSpec = serde_yaml::from_str(yaml).unwrap();
let compiled = CompiledSchedule::compile(&spec).unwrap();
SchedulerState::new(&compiled)
}
#[test]
fn dispatches_when_idle() {
let mut s = state("cron: \"* * * * *\"");
assert_eq!(s.on_tick(false), TickAction::Dispatch);
}
#[test]
fn skip_policy_drops_overlapping_tick() {
let mut s = state("cron: \"* * * * *\"\noverlap_policy: skip");
assert_eq!(s.on_tick(true), TickAction::Skip);
}
#[test]
fn queue_policy_buffers_and_dispatches_on_completion() {
let mut s = state("cron: \"* * * * *\"\noverlap_policy: queue");
assert_eq!(s.on_tick(true), TickAction::Queue);
assert_eq!(
s.on_run_finished(RunOutcome::Success),
AfterRun::Continue {
dispatch_pending: true
}
);
assert_eq!(
s.on_run_finished(RunOutcome::Success),
AfterRun::Continue {
dispatch_pending: false
}
);
}
#[test]
fn forbid_policy_aborts_on_overlap() {
let mut s = state("cron: \"* * * * *\"\noverlap_policy: forbid");
assert_eq!(s.on_tick(true), TickAction::ForbidAbort);
}
#[test]
fn max_runs_counts_successes_only() {
let mut s = state("cron: \"* * * * *\"\nmax_runs: 2");
assert_eq!(
s.on_run_finished(RunOutcome::Failure),
AfterRun::Continue {
dispatch_pending: false
}
);
assert_eq!(
s.on_run_finished(RunOutcome::Success),
AfterRun::Continue {
dispatch_pending: false
}
);
assert_eq!(s.on_run_finished(RunOutcome::Success), AfterRun::ExitOk);
}
#[test]
fn on_failure_stop_exits_on_first_failure() {
let mut s = state("cron: \"* * * * *\"\non_failure: stop");
assert_eq!(
s.on_run_finished(RunOutcome::Failure),
AfterRun::ExitFailure { consecutive: 1 }
);
}
#[test]
fn max_consecutive_failures_trips_and_success_resets() {
let mut s = state("cron: \"* * * * *\"\nmax_consecutive_failures: 3");
assert_eq!(
s.on_run_finished(RunOutcome::Failure),
AfterRun::Continue {
dispatch_pending: false
}
);
assert_eq!(
s.on_run_finished(RunOutcome::Success),
AfterRun::Continue {
dispatch_pending: false
}
);
assert_eq!(
s.on_run_finished(RunOutcome::Failure),
AfterRun::Continue {
dispatch_pending: false
}
);
assert_eq!(
s.on_run_finished(RunOutcome::Failure),
AfterRun::Continue {
dispatch_pending: false
}
);
assert_eq!(
s.on_run_finished(RunOutcome::Failure),
AfterRun::ExitFailure { consecutive: 3 }
);
}
}