//! Job runner: starts a handler for a job kind, persists its log and result. use std::collections::HashMap; use std::sync::Arc; use async_trait::async_trait; use chrono::Utc; use domain::jobs::{JobKind, JobRun, JobStatus}; use domain::ports::JobRunRepository; use domain::DomainError; use uuid::Uuid; /// Sink for log lines of a running job. #[async_trait] pub trait JobLog: Send + Sync { async fn line(&self, text: &str); } #[async_trait] pub trait JobHandler: Send + Sync { async fn run(&self, params: Option, log: &dyn JobLog) -> Result<(), String>; } pub struct JobRunner { runs: Arc, handlers: HashMap>, } struct RepoLog { runs: Arc, id: Uuid, } #[async_trait] impl JobLog for RepoLog { async fn line(&self, text: &str) { if let Err(e) = self.runs.append_log(self.id, text).await { tracing_line(&format!("failed to append job log: {e}")); } } } fn tracing_line(msg: &str) { eprintln!("{msg}"); } impl JobRunner { pub fn new(runs: Arc) -> Self { Self { runs, handlers: HashMap::new(), } } pub fn register(mut self, kind: JobKind, handler: Arc) -> Self { self.handlers.insert(kind, handler); self } pub fn kinds(&self) -> Vec { let mut k: Vec<_> = self.handlers.keys().copied().collect(); k.sort_by_key(|k| k.as_str()); k } async fn begin( &self, kind: JobKind, params: Option, triggered_by: &str, ) -> Result<(JobRun, Arc), DomainError> { let handler = self .handlers .get(&kind) .cloned() .ok_or(DomainError::NotFound)?; if self.runs.find_running(kind).await?.is_some() { return Err(DomainError::Conflict(format!( "{} is already running", kind.as_str() ))); } let run = JobRun { id: Uuid::new_v4(), kind, params, status: JobStatus::Running, started_at: Utc::now(), finished_at: None, log: String::new(), triggered_by: triggered_by.into(), }; self.runs.insert(&run).await?; Ok((run, handler)) } async fn execute(runs: Arc, handler: Arc, run: &JobRun) { let log = RepoLog { runs: runs.clone(), id: run.id, }; let status = match handler.run(run.params.clone(), &log).await { Ok(()) => JobStatus::Success, Err(e) => { log.line(&format!("ERROR: {e}")).await; JobStatus::Failed } }; if let Err(e) = runs.finish(run.id, status).await { tracing_line(&format!("failed to finish job: {e}")); } } /// Start a job in the background. Fails with `Conflict` if the kind is already running. pub async fn start( &self, kind: JobKind, params: Option, triggered_by: &str, ) -> Result { let (run, handler) = self.begin(kind, params, triggered_by).await?; let (runs, run_clone) = (self.runs.clone(), run.clone()); tokio::spawn(async move { Self::execute(runs, handler, &run_clone).await }); Ok(run) } /// Run a job and wait for it to finish (used by tests and the scheduler). pub async fn run_and_wait( &self, kind: JobKind, params: Option, triggered_by: &str, ) -> Result { let (run, handler) = self.begin(kind, params, triggered_by).await?; Self::execute(self.runs.clone(), handler, &run).await; self.get(run.id).await } pub async fn get(&self, id: Uuid) -> Result { self.runs.get(id).await?.ok_or(DomainError::NotFound) } pub async fn list(&self, limit: u32) -> Result, DomainError> { self.runs.list(limit).await } pub async fn last_finished(&self, kind: JobKind) -> Result, DomainError> { self.runs.last_finished(kind).await } }