use crate::prelude::*;
pub struct JobRunner {
pub semaphore: Arc<Semaphore>,
pub set: RefMut<JoinSet<Result<(), Failure<JobAction>>>>,
pub publisher: Ref<Publisher>,
}
#[injectable]
impl JobRunner {
pub fn new(
semaphore: Arc<Semaphore>,
set: RefMut<JoinSet<Result<(), Failure<JobAction>>>>,
publisher: Ref<Publisher>,
) -> Self {
Self {
semaphore,
set,
publisher,
}
}
pub fn add(&self, jobs: Vec<Job>) {
for job in jobs {
let id = job.get_id();
let semaphore = self.semaphore.clone();
let publisher = self.publisher.clone();
publisher.update(&id, Created);
let mut set = self.set.write().expect("join set to be writeable");
set.spawn(async move {
publisher.update(&id, Queued);
let _permit = semaphore
.acquire()
.await
.expect("Semaphore should be available");
publisher.update(&id, Started);
job.execute().await.map_err(|f| f.with("job", &id))?;
publisher.update(&id, Completed);
Ok(())
});
}
}
pub fn add_without_publish(&self, jobs: Vec<Job>) {
for job in jobs {
let id = job.get_id();
let semaphore = self.semaphore.clone();
let mut set = self.set.write().expect("join set to be writeable");
set.spawn(async move {
let _permit = semaphore
.acquire()
.await
.expect("Semaphore should be available");
job.execute().await.map_err(|f| f.with("job", &id))?;
Ok(())
});
}
}
pub async fn execute(&self) -> Result<(), Failure<JobAction>> {
self.execute_internal(true).await
}
pub async fn execute_without_publish(&self) -> Result<(), Failure<JobAction>> {
self.execute_internal(false).await
}
async fn execute_internal(&self, publish: bool) -> Result<(), Failure<JobAction>> {
if publish {
self.publisher.start("");
}
let mut set = self.set.write().expect("join set to be writeable");
while let Some(result) = set.join_next().await {
let result = match result {
Ok(result) => result,
Err(e) => {
set.abort_all();
set.detach_all();
return Err(Failure::new(JobAction::ExecuteTask, e));
}
};
if let Err(e) = result {
set.abort_all();
set.detach_all();
return Err(e);
}
}
if publish {
self.publisher.finish("");
}
Ok(())
}
}