use azums::jobs::{AttemptsRepo, JobsRepo};
use serde_json::json;
mod common;
use common::setup_db;
#[tokio::test]
async fn worker_crash_mid_job_another_worker_recovers_after_lease_expiry() -> anyhow::Result<()> {
let Some(pool) = setup_db().await else {
return Ok(());
};
let jobs = JobsRepo::new(pool.clone());
let attempts = AttemptsRepo::new(pool.clone());
let job_id = jobs.enqueue_now("default", "fail_me", json!({})).await?;
let job = jobs
.lease_one_job("default", "workerA", 1)
.await?
.expect("leased");
assert_eq!(job.id, job_id);
let _attempt1 = attempts.start_attempt(job_id, "workerA").await?;
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
jobs.reap_expired_locks().await?;
let job2 = jobs
.lease_one_job("default", "workerB", 10)
.await?
.expect("leased by B");
assert_eq!(job2.id, job_id);
let attempt2 = attempts.start_attempt(job_id, "workerB").await?;
attempts
.finish_failed(attempt2.id, 2, "recovered after crash", "TIMEOUT")
.await?;
let status: String = sqlx::query_scalar("SELECT status FROM jobs WHERE id = $1")
.bind(job_id)
.fetch_one(&pool)
.await?;
assert!(
matches!(
status.as_str(),
"running" | "queued" | "dlq" | "failed" | "succeeded"
),
"unexpected status: {status}"
);
let attempt_count: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM job_attempts WHERE job_id = $1")
.bind(job_id)
.fetch_one(&pool)
.await?;
assert_eq!(attempt_count, 2);
Ok(())
}