pub mod error_handling;
pub mod job_group;
pub mod opaque_job;
pub mod trigger;
mod disabled_job;
mod simple_job;
mod job_result;
pub use self::{
disabled_job::DisabledJob, error_handling::HandleError, job_group::JobGroup,
job_result::JobResult, opaque_job::OpaqueJob, trigger::Trigger,
};
use futures::FutureExt;
use non_non_full::NonEmptyVec;
use simple_job::{SimpleJob, SimpleJobBuilder};
use std::ops::ControlFlow;
use std::panic;
use tokio::select;
use self::error_handling::{HandleErrorContext, HandleErrorResult};
use self::trigger::TriggerResult;
use crate::{
StaticStr,
cancellation_token::{CancellationToken, cancel_wait},
error::ErrorChainDisplay,
maybe_send::MaybeSync,
task::TaskGroup,
};
#[derive(bon::Builder, Clone, Debug)]
#[builder(finish_fn(name = "build_internal", vis = ""))]
#[builder(builder_type(doc {
/// Use builder syntax to set the inputs and finish with [`build()`](`JobBuilder::build()`)
/// or [`build_with_default_error_handling()`](`JobBuilder::build_with_default_error_handling()`).
}))]
#[non_exhaustive]
pub struct Job<T, Tr, H> {
#[builder(start_fn, into)]
pub name: StaticStr,
pub tasks: T,
pub trigger: Tr,
pub error_handling: H,
#[builder(required)]
pub cancel_token: Option<CancellationToken>,
}
impl<Tr, H> Job<(), Tr, H> {
#[cfg_attr(all(feature = "source-http", feature = "action-html"), doc = "```")]
#[cfg_attr(
not(all(feature = "source-http", feature = "action-html")),
doc = " ```ignore"
)]
#[doc = "```"]
pub fn builder_simple<S, A>(name: impl Into<StaticStr>) -> SimpleJobBuilder<S, A, Tr, H> {
SimpleJob::builder(name)
}
}
impl<T, Tr, H> Job<T, Tr, H>
where
T: TaskGroup,
Tr: Trigger + MaybeSync,
H: HandleError<Tr>,
{
#[expect(clippy::same_name_method, reason = "can't think of a better name")] pub async fn run(&mut self) -> JobResult {
match panic::AssertUnwindSafe(self.run_inner())
.catch_unwind()
.await
{
Ok(job_result) => job_result,
Err(panic_payload) => JobResult::Panicked {
payload: panic_payload,
},
}
}
async fn run_inner(&mut self) -> JobResult {
tracing::info!("Starting job {}", self.name);
loop {
let results = self.tasks.run_concurrently().await;
#[expect(clippy::manual_ok_err, reason = "false positive")]
let errors = results
.into_iter()
.filter_map(|r| {
tracing::trace!("Task result: {r:?}");
match r {
Ok(()) => None,
Err(e) => Some(e),
}
})
.collect::<Vec<_>>();
if let Some(errors) = NonEmptyVec::new(errors) {
let cx = HandleErrorContext {
job_name: &self.name,
job_trigger: &self.trigger,
cancel_token: self.cancel_token.as_mut(),
};
match self.error_handling.handle_errors(errors, cx).await {
HandleErrorResult::ResumeJob {
wait_for_trigger: wait_on_the_trigger,
} => {
if !wait_on_the_trigger {
continue;
}
}
HandleErrorResult::StopWithErrors(e) => return JobResult::Err(e),
HandleErrorResult::ErrWhileHandling {
err,
original_errors,
} => {
tracing::error!(
"An error occured while handling other errors! Stopping the job and returning the original errors.\nDetails: {err}",
);
return JobResult::Err(original_errors);
}
}
}
match wait_for_trigger(&mut self.trigger, self.cancel_token.as_mut(), &self.name).await
{
ControlFlow::Continue(()) => (),
ControlFlow::Break(res) => return res,
}
}
}
}
impl<T, Tr, H> OpaqueJob for Job<T, Tr, H>
where
T: TaskGroup,
Tr: Trigger + MaybeSync,
H: HandleError<Tr>,
{
async fn run(&mut self) -> JobResult {
Job::run(self).await
}
fn name(&self) -> Option<&str> {
Some(&self.name)
}
}
impl<T, Tr, S: job_builder::State> JobBuilder<T, Tr, error_handling::ExponentialBackoff, S>
where
T: TaskGroup,
{
pub fn build_with_default_error_handling(self) -> Job<T, Tr, error_handling::ExponentialBackoff>
where
S::CancelToken: job_builder::IsSet,
S::ErrorHandling: job_builder::IsUnset,
S::Trigger: job_builder::IsSet,
S::Tasks: job_builder::IsSet,
{
let this = self.error_handling(error_handling::ExponentialBackoff::default());
this.build()
}
}
impl<T, Tr, H, S: job_builder::State> JobBuilder<T, Tr, H, S>
where
T: TaskGroup,
{
pub fn build(self) -> Job<T, Tr, H>
where
S: job_builder::IsComplete,
S::CancelToken: job_builder::IsSet,
S::ErrorHandling: job_builder::IsSet,
S::Trigger: job_builder::IsSet,
S::Tasks: job_builder::IsSet,
{
let mut job = self.build_internal();
if let Some(token) = &job.cancel_token {
job.tasks.set_cancel_token(token.clone());
}
job
}
}
async fn wait_for_trigger<Tr>(
mut trigger: Tr,
cancel_token: Option<&mut CancellationToken>,
job_name: &str,
) -> ControlFlow<JobResult>
where
Tr: Trigger,
{
select! {
trigger_res = trigger.wait() => {
match trigger_res {
Ok(TriggerResult::Resume) => ControlFlow::Continue(()),
Ok(TriggerResult::Stop) => ControlFlow::Break(JobResult::Ok),
Err(e) => ControlFlow::Break(JobResult::TriggerFailed(e.into())),
}
},
() = cancel_wait(cancel_token) => {
tracing::info!("Job {job_name} is shutting down...");
ControlFlow::Break(JobResult::Ok)
}
}
}