Implement Trivy adapter, contain job panics, fail interrupted runs on startup
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
CI / ui (push) Has been cancelled

The Trivy scanner still had unimplemented stubs, which panicked the scan
task in the deployed test instance and left the run in 'running' forever.
JobRunner now runs handlers in their own task and marks a panic as a
failed run; on startup runs left 'running' by a previous process are
marked failed.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Dennis Nemec
2026-09-02 22:58:16 +02:00
parent 872f4373ff
commit 5026ce22de
7 changed files with 212 additions and 15 deletions

View File

@ -94,22 +94,41 @@ impl JobRunner {
}
async fn execute(runs: Arc<dyn JobRunRepository>, handler: Arc<dyn JobHandler>, run: &JobRun) {
let log = RepoLog {
let log = Arc::new(RepoLog {
runs: runs.clone(),
id: run.id,
};
let status = match handler.run(run.params.clone(), &log).await {
Ok(()) => JobStatus::Success,
Err(e) => {
});
// Run the handler in its own task so a panic is contained and reported.
let (params, task_log) = (run.params.clone(), log.clone());
let joined = tokio::spawn(async move { handler.run(params, &*task_log).await }).await;
let status = match joined {
Ok(Ok(())) => JobStatus::Success,
Ok(Err(e)) => {
log.line(&format!("ERROR: {e}")).await;
JobStatus::Failed
}
Err(e) => {
log.line(&format!("ERROR: job panicked: {e}")).await;
JobStatus::Failed
}
};
if let Err(e) = runs.finish(run.id, status).await {
tracing_line(&format!("failed to finish job: {e}"));
}
}
/// Mark runs left in `running` state by a previous process as failed. Returns the count.
pub async fn recover(&self) -> Result<usize, DomainError> {
let stale = self.runs.running().await?;
for r in &stale {
self.runs
.append_log(r.id, "ERROR: interrupted by a restart of the service")
.await?;
self.runs.finish(r.id, JobStatus::Failed).await?;
}
Ok(stale.len())
}
/// Start a job in the background. Fails with `Conflict` if the kind is already running.
pub async fn start(
&self,

View File

@ -292,6 +292,16 @@ impl JobRunRepository for MemJobRuns {
.find(|r| r.kind == kind && r.status != JobStatus::Running)
.cloned())
}
async fn running(&self) -> Result<Vec<JobRun>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.filter(|r| r.status == JobStatus::Running)
.cloned()
.collect())
}
}
pub fn smtp() -> SmtpSettings {

View File

@ -7,6 +7,7 @@ use uuid::Uuid;
use crate::jobs::{JobHandler, JobLog, JobRunner};
use crate::test_fakes::MemJobRuns;
use domain::ports::JobRunRepository;
struct Echo;
#[async_trait]
@ -102,3 +103,48 @@ async fn list_returns_newest_first_with_limit() {
assert_eq!(list.len(), 2);
assert!(list[0].started_at >= list[1].started_at);
}
struct Panics;
#[async_trait]
impl JobHandler for Panics {
async fn run(&self, _: Option<String>, log: &dyn JobLog) -> Result<(), String> {
log.line("about to panic").await;
panic!("handler bug");
}
}
#[tokio::test]
async fn panicking_handler_marks_run_failed() {
let runs = Arc::new(MemJobRuns::default());
let r = JobRunner::new(runs.clone()).register(JobKind::Backup, Arc::new(Panics));
let run = r.start(JobKind::Backup, None, "test").await.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
let run = r.get(run.id).await.unwrap();
assert_eq!(run.status, JobStatus::Failed);
assert!(run.log.contains("panicked"), "{}", run.log);
// the kind is free again
assert!(r.start(JobKind::Backup, None, "test").await.is_ok());
}
#[tokio::test]
async fn stale_running_runs_are_failed_on_recovery() {
let runs = Arc::new(MemJobRuns::default());
runs.insert(&domain::jobs::JobRun {
id: Uuid::new_v4(),
kind: JobKind::PackageRefresh,
params: None,
status: JobStatus::Running,
started_at: chrono::Utc::now(),
finished_at: None,
log: String::new(),
triggered_by: "old process".into(),
})
.await
.unwrap();
let r = JobRunner::new(runs.clone()).register(JobKind::PackageRefresh, Arc::new(Echo));
assert_eq!(r.recover().await.unwrap(), 1);
let list = r.list(10).await.unwrap();
assert_eq!(list[0].status, JobStatus::Failed);
assert!(list[0].log.contains("interrupted"));
assert!(r.start(JobKind::PackageRefresh, None, "test").await.is_ok());
}