use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use crate::application::resources::Resources;
use crate::application::{EngineError, EngineResult};
pub struct JobsRuntime {
shutdown: CancellationToken,
worker: JoinHandle<Result<(), crate::jobs::WorkerError>>,
scheduler: Option<JoinHandle<Result<(), crate::jobs::SchedulerError>>>,
}
impl JobsRuntime {
pub async fn shutdown(self) -> EngineResult<()> {
self.shutdown.cancel();
if let Err(join_err) = self.worker.await {
return Err(EngineError::Shutdown {
subsystem: "jobs",
source: crate::Error::Job(format!("worker task panicked: {join_err}")),
});
}
if let Some(scheduler) = self.scheduler
&& let Err(join_err) = scheduler.await
{
return Err(EngineError::Shutdown {
subsystem: "jobs",
source: crate::Error::Job(format!("scheduler task panicked: {join_err}")),
});
}
Ok(())
}
}
pub(super) async fn start_jobs(
registry: Option<crate::jobs::Registry>,
worker_config: Option<crate::jobs::WorkerConfig>,
scheduler: Option<crate::jobs::Scheduler>,
resources: &mut Resources,
) -> EngineResult<Option<JobsRuntime>> {
let Some(registry) = registry else {
return Ok(None);
};
let db = resources.db().ok_or_else(|| EngineError::Startup {
subsystem: "jobs",
stage: "connect",
source: crate::Error::Config(
"the jobs subsystem requires the database subsystem to be enabled and connected"
.to_string(),
),
})?;
let pool = db.sqlx().clone();
let jobs = crate::jobs::Jobs::new(pool.clone());
jobs.migrate().await.map_err(|e| EngineError::Startup {
subsystem: "jobs",
stage: "migrate",
source: crate::Error::Job(e.to_string()),
})?;
resources.set_jobs(jobs);
let shutdown = CancellationToken::new();
let worker_config = worker_config.unwrap_or_default();
let worker = crate::jobs::Worker::builder(pool, registry)
.config(worker_config)
.build();
let worker_shutdown = shutdown.clone();
let worker_handle = tokio::spawn(async move { worker.run(worker_shutdown).await });
let scheduler_handle = if let Some(scheduler) = scheduler {
let scheduler_shutdown = shutdown.clone();
Some(tokio::spawn(async move {
scheduler.run(scheduler_shutdown).await
}))
} else {
None
};
Ok(Some(JobsRuntime {
shutdown,
worker: worker_handle,
scheduler: scheduler_handle,
}))
}