use dashmap::DashMap;
use snafu::{ResultExt, Snafu};
use std::sync::LazyLock;
use tibba_error::Error as BaseError;
pub use tokio_cron_scheduler::Job;
use tokio_cron_scheduler::{JobScheduler, JobSchedulerError};
use tracing::info;
const LOG_TARGET: &str = "tibba:scheduler";
type Result<T> = std::result::Result<T, BaseError>;
#[derive(Debug, Snafu)]
enum Error {
#[snafu(display("create scheduler failed: {source}"))]
Create { source: JobSchedulerError },
#[snafu(display("add job {name} failed: {source}"))]
AddJob {
name: String,
source: JobSchedulerError,
},
#[snafu(display("start scheduler failed: {source}"))]
Start { source: JobSchedulerError },
}
impl From<Error> for BaseError {
fn from(val: Error) -> Self {
let err = match val {
Error::Create { source } => BaseError::new(source),
Error::AddJob { name, source } => BaseError::new(source).with_sub_category(name),
Error::Start { source } => BaseError::new(source),
};
err.with_category("scheduler")
}
}
static JOB_TASKS: LazyLock<DashMap<String, Job>> = LazyLock::new(DashMap::new);
pub fn register_job_task(name: impl Into<String>, job: Job) {
JOB_TASKS.insert(name.into(), job);
}
pub async fn run_scheduler_jobs() -> Result<JobScheduler> {
let scheduler = JobScheduler::new().await.context(CreateSnafu)?;
for item in JOB_TASKS.iter() {
let (name, job) = item.pair();
scheduler
.add(job.clone())
.await
.context(AddJobSnafu { name: name.clone() })?;
info!(target: LOG_TARGET, name, "add job success");
}
scheduler.shutdown_on_ctrl_c();
scheduler.start().await.context(StartSnafu)?;
info!(target: LOG_TARGET, "scheduler started");
Ok(scheduler)
}