use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio_cron_scheduler::{Job, JobScheduler};
use crate::response::ServiceResult;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum JobType {
Periodic,
Exact,
PeriodicAndExact,
}
#[derive(Debug, Clone)]
pub struct TimeOfDay {
pub hour: u8,
pub minute: u8,
pub second: u8
}
impl TimeOfDay {
pub fn new(hour: u8, minute: u8, second: u8) -> Result<Self, String> {
if hour > 23 { return Err("hour >23".into()); }
if minute > 59 { return Err("minute >59".into()); }
if second > 59 { return Err("second >59".into()); }
Ok(Self { hour, minute, second })
}
}
#[async_trait::async_trait]
pub trait ServiceJob: Send + Sync {
fn name(&self) -> &str;
fn job_type(&self) -> JobType;
fn period(&self) -> Option<Duration>;
fn schedule(&self) -> Option<TimeOfDay>;
fn is_critical(&self) -> bool { false }
async fn run(&self) -> anyhow::Result<ServiceResult<serde_json::Value>>;
async fn stop_gracefully(&self) -> anyhow::Result<()> { Ok(()) }
}
#[derive(Clone)]
struct JobEntry {
job: Arc<dyn ServiceJob>,
runs: Arc<Mutex<usize>>,
failures: Arc<Mutex<usize>>,
successes: Arc<Mutex<usize>>,
}
pub struct JobRegistry {
jobs: Vec<JobEntry>,
scheduler: Option<JobScheduler>,
}
impl JobRegistry {
pub fn new() -> Self { Self { jobs: vec![], scheduler: None } }
pub fn get_runs(&self, name: &str) -> Option<usize> {
for entry in &self.jobs {
if entry.job.name() == name {
return Some(entry.runs.lock().unwrap().clone());
}
}
None
}
pub fn get_failures(&self, name: &str) -> Option<usize> {
for entry in &self.jobs {
if entry.job.name() == name {
return Some(entry.failures.lock().unwrap().clone());
}
}
None
}
pub fn get_successes(&self, name: &str) -> Option<usize> {
for entry in &self.jobs {
if entry.job.name() == name {
return Some(entry.successes.lock().unwrap().clone());
}
}
None
}
pub fn add_job<J: ServiceJob + 'static>(&mut self, job: J) {
self.jobs.push(JobEntry {
job: Arc::new(job),
runs: Arc::new(Mutex::new(0)),
failures: Arc::new(Mutex::new(0)),
successes: Arc::new(Mutex::new(0)),
});
}
pub fn job_count(&self) -> usize { self.jobs.len() }
pub async fn start(&mut self) -> anyhow::Result<()> {
let sched = JobScheduler::new().await?;
tracing::info!(job_count = self.jobs.len(), "Starting job scheduler");
for entry in &self.jobs {
let job_clone = entry.job.clone();
match job_clone.job_type() {
JobType::Periodic => {
let period = job_clone.period().unwrap_or(Duration::from_secs(60));
tracing::info!(job = %job_clone.name(), schedule = "periodic", interval_secs = period.as_secs(), "Registered job");
tokio::spawn({
let job_clone = job_clone.clone();
async move {
let mut interval = tokio::time::interval(period);
loop {
interval.tick().await;
if let Err(e) = job_clone.run().await {
tracing::error!(job = %job_clone.name(), error = %e, "Periodic job failed");
}
}
}
});
}
JobType::PeriodicAndExact => {
let schedule_clone = job_clone.clone();
let period_clone = job_clone.clone();
if let Some(tod) = schedule_clone.schedule() {
let job_name = schedule_clone.name().to_string();
let cron = format!("{} {} {} 1 * *", tod.second, tod.minute, tod.hour);
let j = Job::new_async(cron.as_str(), move |_, _| {
let jc = schedule_clone.clone();
Box::pin(async move {
if let Err(e) = jc.run().await {
tracing::error!(job = %jc.name(), error = %e, "Scheduled job failed");
}
})
})?;
sched.add(j).await?;
tracing::info!(job = %job_name, schedule = "exact", "Registered scheduled job");
}
if let Some(period) = period_clone.period() {
tracing::info!(job = %period_clone.name(), schedule = "periodic", interval_secs = period.as_secs(), "Registered periodic job leg");
let pc = period_clone.clone();
tokio::spawn(async move {
let mut interval = tokio::time::interval(period);
loop {
interval.tick().await;
if let Err(e) = pc.run().await {
tracing::error!(job = %pc.name(), error = %e, "Periodic job leg failed");
}
}
});
}
}
JobType::Exact => {
if let Some(tod) = job_clone.schedule() {
let job_name = job_clone.name().to_string();
let cron = format!("{} {} {} * * *", tod.second, tod.minute, tod.hour);
let j = Job::new_async(cron.as_str(), move |_, _| {
let jc = job_clone.clone();
Box::pin(async move {
if let Err(e) = jc.run().await {
tracing::error!(job = %jc.name(), error = %e, "Scheduled job failed");
}
})
})?;
sched.add(j).await?;
tracing::info!(job = %job_name, schedule = "exact", "Registered scheduled job");
}
}
}
}
sched.start().await?;
self.scheduler = Some(sched);
tracing::info!("Job scheduler started");
Ok(())
}
pub async fn stop(&mut self) -> anyhow::Result<()> {
tracing::info!(job_count = self.jobs.len(), "Stopping jobs");
for entry in &self.jobs {
entry.job.stop_gracefully().await?;
}
if let Some(mut s) = self.scheduler.take() {
s.shutdown().await?;
}
tracing::info!("Jobs stopped");
Ok(())
}
}
impl Default for JobRegistry {
fn default() -> Self { Self::new() }
}