use crate::{event::EventId, job::JobId, AsyncJobBoxed, Error};
use chrono::{DateTime, Local};
use cron_lite::Schedule;
use futures::Future;
use std::{
collections::BTreeSet,
fmt::Display,
pin::Pin,
sync::Arc,
time::{Duration, SystemTime},
};
use tokio::sync::RwLock;
use tracing::debug;
use uuid::Uuid;
#[derive(Clone)]
pub struct Task {
pub(crate) id: TaskId,
pub(crate) job: AsyncJobBoxed,
pub(crate) schedule: TaskSchedule,
pub(crate) state: TaskState,
pub(crate) timeout: Option<Duration>,
}
impl Task {
pub fn new<T>(schedule: TaskSchedule, job: T) -> Self
where
T: 'static,
T: FnMut(JobId) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync,
{
Self {
id: TaskId::new(),
job: Arc::new(RwLock::new(Box::new(job))),
schedule,
state: TaskState::default(),
timeout: None,
}
}
pub fn with_id(self, id: impl Into<TaskId>) -> Self {
Self {
id: id.into(),
..self
}
}
pub fn with_schedule(self, schedule: impl Into<TaskSchedule>) -> Self {
Self {
schedule: schedule.into(),
..self
}
}
pub fn with_timeout(self, timeout: impl Into<Duration>) -> Self {
Self {
timeout: Some(timeout.into()),
..self
}
}
#[deprecated(since = "0.4.2", note = "please use `with_id` method instead.")]
pub fn new_with_id<T>(schedule: TaskSchedule, job: T, id: TaskId) -> Self
where
T: 'static,
T: FnMut(JobId) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync,
{
Self {
id,
job: Arc::new(RwLock::new(Box::new(job))),
schedule,
state: TaskState::default(),
timeout: None,
}
}
pub fn id(&self) -> TaskId {
self.id.clone()
}
pub fn schedule(&self) -> TaskSchedule {
self.schedule.clone()
}
pub fn timeout(&self) -> Option<Duration> {
self.timeout
}
pub fn status(&self) -> TaskStatus {
self.state.status()
}
pub fn statistics(&self) -> TaskStatistics {
self.state.statistics.clone()
}
}
impl std::fmt::Debug for Task {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Task")
.field("id", &self.id)
.field("schedule", &self.schedule)
.field("state", &self.state)
.field("timeout", &self.timeout)
.finish()
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct TaskId {
pub(crate) id: String,
}
impl TaskId {
pub fn new() -> Self {
Self {
id: Uuid::new_v4().into(),
}
}
}
impl Default for TaskId {
fn default() -> Self {
Self::new()
}
}
impl From<&str> for TaskId {
fn from(value: &str) -> Self {
Self { id: value.into() }
}
}
impl From<String> for TaskId {
fn from(value: String) -> Self {
Self { id: value }
}
}
impl From<&String> for TaskId {
fn from(value: &String) -> Self {
Self {
id: value.to_owned(),
}
}
}
impl From<Uuid> for TaskId {
fn from(value: Uuid) -> Self {
Self { id: value.into() }
}
}
impl From<&Uuid> for TaskId {
fn from(value: &Uuid) -> Self {
Self {
id: value.to_string(),
}
}
}
impl From<TaskId> for String {
fn from(value: TaskId) -> Self {
value.id
}
}
impl From<&TaskId> for String {
fn from(value: &TaskId) -> Self {
value.id.to_owned()
}
}
impl Display for TaskId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.id)
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum TaskSchedule {
Once,
OnceDelayed(Duration),
Interval(Duration),
IntervalDelayed(Duration, Duration),
Cron(CronSchedule, CronOpts),
}
impl TaskSchedule {
pub(crate) fn initial_run_time(&self) -> SystemTime {
match self {
TaskSchedule::Once => SystemTime::now(),
TaskSchedule::OnceDelayed(delay) => SystemTime::now().checked_add(*delay).unwrap(),
TaskSchedule::Interval(_interval) => SystemTime::now(),
TaskSchedule::IntervalDelayed(_interval, delay) => {
SystemTime::now().checked_add(*delay).unwrap()
}
TaskSchedule::Cron(schedule, opts) => {
if opts.at_start {
SystemTime::now()
} else {
schedule.upcoming()
}
}
}
}
pub(crate) fn after_start_run_time(&self) -> Option<SystemTime> {
match self {
TaskSchedule::Cron(schedule, opts) => {
if opts.concurrent {
Some(schedule.upcoming())
} else {
None
}
}
TaskSchedule::Once => None,
TaskSchedule::OnceDelayed(_) => None,
TaskSchedule::Interval(_) => None,
TaskSchedule::IntervalDelayed(_, _) => None,
}
}
pub(crate) fn after_finish_run_time(&self) -> Option<SystemTime> {
match self {
TaskSchedule::Interval(interval) => Some(SystemTime::now().checked_add(*interval)?),
TaskSchedule::IntervalDelayed(interval, _delay) => {
Some(SystemTime::now().checked_add(*interval)?)
}
TaskSchedule::Once => None,
TaskSchedule::OnceDelayed(_) => None,
TaskSchedule::Cron(schedule, opts) => {
if opts.concurrent {
None
} else {
Some(schedule.upcoming())
}
}
}
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct CronOpts {
pub at_start: bool,
pub concurrent: bool,
}
#[derive(Default, Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub enum TaskStatus {
#[default]
New,
Waiting,
Scheduled,
Running,
Finished,
}
#[derive(Default, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct TaskStatistics {
pub waiting: usize,
pub scheduled: usize,
pub running: usize,
pub completed: usize,
pub canceled: usize,
pub timeouts: usize,
pub errors: usize,
}
#[derive(Default, Clone)]
pub(crate) struct TaskState {
statistics: TaskStatistics,
scheduled_jobs: BTreeSet<JobId>,
running_jobs: BTreeSet<JobId>,
last_finished_at: Option<SystemTime>,
}
impl std::fmt::Debug for TaskState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let last_finished_at = if let Some(last_finished_at) = self.last_finished_at {
format!("{}", DateTime::<Local>::from(last_finished_at))
} else {
"None".to_string()
};
f.debug_struct("TaskState")
.field("waiting", &self.statistics.waiting)
.field("scheduled", &self.statistics.scheduled)
.field("running", &self.statistics.running)
.field("completed", &self.statistics.completed)
.field("canceled", &self.statistics.canceled)
.field("timeouts", &self.statistics.timeouts)
.field("errors", &self.statistics.errors)
.field("scheduled_jobs", &self.scheduled_jobs)
.field("running_jobs", &self.running_jobs)
.field("last_finished_at", &last_finished_at)
.finish()
}
}
impl TaskState {
pub(crate) fn status(&self) -> TaskStatus {
if self.statistics.running > 0 {
TaskStatus::Running
} else if self.statistics.scheduled > 0 {
TaskStatus::Scheduled
} else if self.statistics.waiting > 0 {
TaskStatus::Waiting
} else if self.statistics.completed > 0
|| self.statistics.canceled > 0
|| self.statistics.timeouts > 0
|| self.statistics.errors > 0
{
TaskStatus::Finished
} else {
TaskStatus::New
}
}
pub(crate) fn task_enqueued(&mut self) -> &Self {
self.statistics.waiting += 1;
debug!(status = ?self.status(), "task enqueued");
self
}
pub(crate) fn job_scheduled(&mut self, id: JobId) -> &Self {
self.statistics.waiting -= 1;
self.statistics.scheduled += 1;
self.scheduled_jobs.insert(id);
debug!(status = ?self.status(), "job scheduled");
self
}
pub(crate) fn job_started(&mut self, id: JobId) -> &Self {
self.statistics.scheduled -= 1;
self.statistics.running += 1;
self.scheduled_jobs.remove(&id);
self.running_jobs.insert(id);
debug!(status = ?self.status(), "job started");
self
}
pub(crate) fn job_completed(&mut self, id: &JobId) -> &Self {
self.statistics.running -= 1;
self.statistics.completed += 1;
self.running_jobs.remove(id);
self.last_finished_at = Some(SystemTime::now());
debug!(status = ?self.status(), "job completed");
self
}
pub(crate) fn job_canceled(&mut self, id: &JobId) -> &Self {
self.statistics.canceled += 1;
if self.running_jobs.remove(id) {
self.statistics.running -= 1;
} else {
self.scheduled_jobs.remove(id);
self.statistics.scheduled -= 1;
}
self.last_finished_at = Some(SystemTime::now());
debug!(status = ?self.status(), "job canceled");
self
}
pub(crate) fn job_timeout(&mut self, id: &JobId) -> &Self {
self.statistics.running -= 1;
self.statistics.timeouts += 1;
self.running_jobs.remove(id);
self.last_finished_at = Some(SystemTime::now());
debug!(status = ?self.status(), "job killed because timeout");
self
}
pub(crate) fn job_error(&mut self, id: &JobId) -> &Self {
self.statistics.running -= 1;
self.statistics.errors += 1;
self.running_jobs.remove(id);
self.last_finished_at = Some(SystemTime::now());
debug!(status = ?self.status(), "job finished with error");
self
}
pub(crate) fn is_task_finished(&self) -> bool {
self.statistics.waiting == 0
&& self.statistics.scheduled == 0
&& self.statistics.running == 0
&& (self.statistics.completed > 0
|| self.statistics.canceled > 0
|| self.statistics.timeouts > 0
|| self.statistics.errors > 0)
&& self.scheduled_jobs.is_empty()
&& self.running_jobs.is_empty()
}
pub(crate) fn last_finished_at(&self) -> Option<SystemTime> {
self.last_finished_at
}
pub(crate) fn jobs(&self) -> BTreeSet<JobId> {
let mut jobs = BTreeSet::new();
jobs.extend(self.scheduled_jobs.clone());
jobs.extend(self.running_jobs.clone());
jobs
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct CronSchedule {
schedule: Box<Schedule>,
}
impl Display for CronSchedule {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.schedule)
}
}
impl CronSchedule {
fn upcoming(&self) -> SystemTime {
let next: SystemTime = self.schedule.upcoming(&Local::now()).unwrap().into();
next
}
}
impl TryFrom<String> for CronSchedule {
type Error = Error;
fn try_from(value: String) -> Result<Self, Self::Error> {
let schedule = Schedule::new(value)?;
Ok(Self {
schedule: Box::new(schedule),
})
}
}
impl TryFrom<&str> for CronSchedule {
type Error = Error;
fn try_from(value: &str) -> Result<Self, Self::Error> {
Self::try_from(value.to_string())
}
}
impl TryFrom<&String> for CronSchedule {
type Error = Error;
fn try_from(value: &String) -> Result<Self, Self::Error> {
Self::try_from(value.to_string())
}
}
impl From<EventId> for TaskId {
fn from(value: EventId) -> Self {
Self { id: value.id }
}
}
#[cfg(test)]
mod test {
use super::*;
use crate::Result;
#[test]
fn cron_with_seconds() {
let cron: Result<CronSchedule> = " * * * * * *".try_into();
assert!(cron.is_ok());
let cron: Result<CronSchedule> = " 5 * * * * * ".try_into();
assert!(cron.is_ok());
let cron: Result<CronSchedule> = " */5 * * * * * ".try_into();
assert!(cron.is_ok());
}
#[test]
fn cron_without_seconds() {
let cron: Result<CronSchedule> = " * * * * *".try_into();
assert!(cron.is_ok());
let cron: Result<CronSchedule> = " 12 * * * * ".try_into();
assert!(cron.is_ok());
let cron: Result<CronSchedule> = " */10 3 * * * ".try_into();
assert!(cron.is_ok());
}
#[test]
fn wrong_cron_expression() {
let cron: Result<CronSchedule> = "* * * * ".try_into();
assert!(cron.is_err());
let cron: Result<CronSchedule> = "* 24 * * * ".try_into();
assert!(cron.is_err());
let cron: Result<CronSchedule> = "* * 0,32 * * ".try_into();
assert!(cron.is_err());
let cron: Result<CronSchedule> = "* * * 13 * ".try_into();
assert!(cron.is_err());
}
#[test]
fn task_state_finished_status_completed() {
let job = JobId::new("job");
let mut state = TaskState::default();
assert_eq!(state.status(), TaskStatus::New);
state.task_enqueued();
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Waiting);
state.job_scheduled(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Scheduled);
state.job_started(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Running);
state.job_completed(&job);
assert!(state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Finished);
}
#[test]
fn task_state_finished_status_canceled_started() {
let job = JobId::new("job");
let mut state = TaskState::default();
assert_eq!(state.status(), TaskStatus::New);
state.task_enqueued();
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Waiting);
state.job_scheduled(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Scheduled);
state.job_started(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Running);
state.job_canceled(&job);
assert!(state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Finished);
}
#[test]
fn task_state_finished_status_canceled_scheduled() {
let job = JobId::new("job");
let mut state = TaskState::default();
assert_eq!(state.status(), TaskStatus::New);
state.task_enqueued();
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Waiting);
state.job_scheduled(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Scheduled);
state.job_canceled(&job);
assert!(state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Finished);
}
#[test]
fn task_state_finished_status_timeout() {
let job = JobId::new("job");
let mut state = TaskState::default();
assert_eq!(state.status(), TaskStatus::New);
state.task_enqueued();
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Waiting);
state.job_scheduled(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Scheduled);
state.job_started(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Running);
state.job_timeout(&job);
assert!(state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Finished);
}
#[test]
fn task_state_finished_status_error() {
let job = JobId::new("job");
let mut state = TaskState::default();
assert_eq!(state.status(), TaskStatus::New);
state.task_enqueued();
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Waiting);
state.job_scheduled(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Scheduled);
state.job_started(job.clone());
assert!(!state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Running);
state.job_error(&job);
assert!(state.is_task_finished());
assert_eq!(state.status(), TaskStatus::Finished);
}
#[test]
fn task_state_transition() {
let job1 = JobId::new("task1");
let job2 = JobId::new("task2");
let job3 = JobId::new("task3");
let job4 = JobId::new("task4");
let job5 = JobId::new("task5");
let mut state = TaskState::default();
assert_eq!(state.status(), TaskStatus::New);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(format!("{state:?}"), String::from("TaskState { waiting: 0, scheduled: 0, running: 0, completed: 0, canceled: 0, timeouts: 0, errors: 0, scheduled_jobs: {}, running_jobs: {}, last_finished_at: \"None\" }"));
assert_eq!(state.statistics, TaskStatistics::default());
state.task_enqueued();
assert_eq!(state.status(), TaskStatus::Waiting);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 1,
..Default::default()
}
);
state.task_enqueued();
assert_eq!(state.status(), TaskStatus::Waiting);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 2,
..Default::default()
}
);
state.task_enqueued();
assert_eq!(state.status(), TaskStatus::Waiting);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 3,
..Default::default()
}
);
state.task_enqueued();
assert_eq!(state.status(), TaskStatus::Waiting);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 4,
..Default::default()
}
);
state.task_enqueued();
assert_eq!(state.status(), TaskStatus::Waiting);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 5,
..Default::default()
}
);
state.job_scheduled(job1.clone());
assert_eq!(state.status(), TaskStatus::Scheduled);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(format!("{state:?}"), format!("TaskState {{ waiting: 4, scheduled: 1, running: 0, completed: 0, canceled: 0, timeouts: 0, errors: 0, scheduled_jobs: {{JobId {{ task_id: TaskId {{ id: \"task1\" }}, id: {} }}}}, running_jobs: {{}}, last_finished_at: \"None\" }}", job1.id.to_string()));
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 4,
scheduled: 1,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone()]);
assert_eq!(jobs, expected);
state.job_started(job1.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 4,
scheduled: 0,
running: 1,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone()]);
assert_eq!(jobs, expected);
state.job_scheduled(job5.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 3,
scheduled: 1,
running: 1,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_started(job5.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 3,
scheduled: 0,
running: 2,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_scheduled(job2.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 2,
scheduled: 1,
running: 2,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone(), job2.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_scheduled(job3.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 1,
scheduled: 2,
running: 2,
..Default::default()
}
);
let jobs = state.jobs();
let expected =
BTreeSet::<JobId>::from([job1.clone(), job2.clone(), job3.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_started(job3.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 1,
scheduled: 1,
running: 3,
..Default::default()
}
);
let jobs = state.jobs();
let expected =
BTreeSet::<JobId>::from([job1.clone(), job2.clone(), job3.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_scheduled(job4.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 2,
running: 3,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([
job1.clone(),
job2.clone(),
job3.clone(),
job4.clone(),
job5.clone(),
]);
assert_eq!(jobs, expected);
state.job_started(job4.clone());
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_none());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 1,
running: 4,
..Default::default()
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([
job1.clone(),
job2.clone(),
job3.clone(),
job4.clone(),
job5.clone(),
]);
assert_eq!(jobs, expected);
state.job_canceled(&job2);
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_some());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 0,
running: 4,
canceled: 1,
..Default::default()
}
);
let jobs = state.jobs();
let expected =
BTreeSet::<JobId>::from([job1.clone(), job3.clone(), job4.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_timeout(&job3);
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_some());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 0,
running: 3,
canceled: 1,
timeouts: 1,
errors: 0,
completed: 0
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone(), job4.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_error(&job4);
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_some());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 0,
running: 2,
canceled: 1,
timeouts: 1,
errors: 1,
completed: 0
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone(), job5.clone()]);
assert_eq!(jobs, expected);
state.job_canceled(&job5);
assert_eq!(state.status(), TaskStatus::Running);
assert!(!state.is_task_finished());
assert!(state.last_finished_at().is_some());
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 0,
running: 1,
canceled: 2,
timeouts: 1,
errors: 1,
completed: 0
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([job1.clone()]);
assert_eq!(jobs, expected);
state.job_completed(&job1);
assert_eq!(state.status(), TaskStatus::Finished);
assert!(state.is_task_finished());
assert!(state.last_finished_at().is_some());
assert_eq!(format!("{state:?}"), format!("TaskState {{ waiting: 0, scheduled: 0, running: 0, completed: 1, canceled: 2, timeouts: 1, errors: 1, scheduled_jobs: {{}}, running_jobs: {{}}, last_finished_at: \"{}\" }}", DateTime::<Local>::from(state.last_finished_at().unwrap())));
assert_eq!(
state.statistics,
TaskStatistics {
waiting: 0,
scheduled: 0,
running: 0,
canceled: 2,
timeouts: 1,
errors: 1,
completed: 1
}
);
let jobs = state.jobs();
let expected = BTreeSet::<JobId>::from([]);
assert_eq!(jobs, expected);
}
#[test]
fn cron_opts_default() {
assert_eq!(
CronOpts::default(),
CronOpts {
at_start: false,
concurrent: false
}
);
}
#[test]
fn type_convertors() {
let uuid_id = Uuid::new_v4();
let str_id = uuid_id.to_string();
assert_eq!(
TaskId::from(String::from("TASK_ID")).id,
String::from("TASK_ID")
);
assert_eq!(
TaskId::from(&String::from("TASK_ID")).id,
String::from("TASK_ID")
);
assert_eq!(TaskId::from("TASK_ID").id, String::from("TASK_ID"));
assert_eq!(TaskId::from(uuid_id).id, str_id);
assert_eq!(TaskId::from(&uuid_id).id, str_id);
assert_eq!(String::from(TaskId::from(uuid_id)), str_id);
assert_eq!(String::from(&TaskId::from(uuid_id)), str_id);
}
#[test]
fn constructors() {
let task = Task::new(TaskSchedule::OnceDelayed(Duration::from_secs(1)), |_id| {
Box::pin(async move {})
});
assert_eq!(
task.clone().with_id("TEST").id().to_string(),
String::from("TEST")
);
assert_eq!(
task.clone().with_schedule(TaskSchedule::Once).schedule(),
TaskSchedule::Once
);
assert_eq!(
task.clone().with_timeout(Duration::from_secs(10)).timeout(),
Some(Duration::from_secs(10))
);
assert_eq!(task.clone().timeout(), None);
assert_eq!(task.status(), TaskStatus::New);
let id = Uuid::new_v4();
#[allow(deprecated)]
let task = Task::new_with_id(
TaskSchedule::OnceDelayed(Duration::from_secs(1)),
|_id| Box::pin(async move {}),
id.into(),
);
assert_eq!(task.id().to_string(), id.to_string());
assert_eq!(task.status(), TaskStatus::New);
assert_eq!(
CronSchedule::try_from("1 2 3 4 ?").unwrap(),
CronSchedule::try_from(String::from("0 1 2 3 4 ? *")).unwrap()
);
assert_eq!(
CronSchedule::try_from("1 2 ? 4 5").unwrap(),
CronSchedule::try_from(&String::from("0 1 2 ? 4 5 *")).unwrap()
);
}
#[test]
fn task_status_reflects_real_task_state() {
let mut task = Task::new(TaskSchedule::Once, |_id| Box::pin(async move {}));
let task_id = task.id();
let job_id = JobId::new(task_id);
assert_eq!(task.status(), TaskStatus::New);
task.state.task_enqueued();
assert_eq!(task.status(), TaskStatus::Waiting);
task.state.job_scheduled(job_id.clone());
assert_eq!(task.status(), TaskStatus::Scheduled);
task.state.job_started(job_id.clone());
assert_eq!(task.status(), TaskStatus::Running);
task.state.job_completed(&job_id);
assert_eq!(task.status(), TaskStatus::Finished);
}
#[test]
fn debug_formatter() {
let task1 = Task::new(TaskSchedule::Once, |_id| Box::pin(async move {})).with_id("TEST");
let task2 = Task::new(TaskSchedule::Once, |_id| Box::pin(async move {}))
.with_id("TEST_WITH_TIMEOUT")
.with_timeout(Duration::from_secs(1));
assert_eq!(format!("{task1:?}"), format!("Task {{ id: TaskId {{ id: \"TEST\" }}, schedule: Once, state: TaskState {{ waiting: 0, scheduled: 0, running: 0, completed: 0, canceled: 0, timeouts: 0, errors: 0, scheduled_jobs: {{}}, running_jobs: {{}}, last_finished_at: \"None\" }}, timeout: None }}"));
assert_eq!(format!("{task2:?}"), format!("Task {{ id: TaskId {{ id: \"TEST_WITH_TIMEOUT\" }}, schedule: Once, state: TaskState {{ waiting: 0, scheduled: 0, running: 0, completed: 0, canceled: 0, timeouts: 0, errors: 0, scheduled_jobs: {{}}, running_jobs: {{}}, last_finished_at: \"None\" }}, timeout: Some(1s) }}"));
assert_eq!(
format!("{}", CronSchedule::try_from("1 2 3 4 ?").unwrap()),
String::from("0 1 2 3 4 ? *")
);
}
}