//! DB-layer contract tests for the scan-job reaper's liveness heartbeat. //! //! Audit Run 22: `reap_stuck` keyed off claim-time `started_at`, which spans the //! full S3 download + scan, so a slow-but-progressing large scan was reaped and //! double-processed, each re-claim inflating `attempts` until a valid file was //! force-retired to `failed` and the entity stranded at HeldForReview. The //! reaper now keys off `heartbeat_at` (falling back to `started_at`). These pin //! that a fresh beat spares a long-running job, a stale beat is still reaped, the //! NULL fallback works, and `bump_heartbeat` refreshes liveness without //! resurrecting a terminal row. use crate::harness::TestHarness; use makenotwork::db::UserId; use makenotwork::db::scan_jobs::{self, ScanTargetKind}; use makenotwork::storage::FileType; use uuid::Uuid; async fn seed_user(h: &TestHarness, name: &str) -> UserId { let hash = makenotwork::auth::hash_password("password123").expect("hash"); sqlx::query_scalar::<_, UserId>( "INSERT INTO users (username, email, password_hash, email_verified) VALUES ($1, $2, $3, true) RETURNING id", ) .bind(name) .bind(format!("{name}@test.com")) .bind(&hash) .fetch_one(&h.db) .await .expect("seed user") } async fn status_of(h: &TestHarness, job: Uuid) -> String { sqlx::query_scalar::<_, String>("SELECT status FROM scan_jobs WHERE id = $1") .bind(job) .fetch_one(&h.db) .await .expect("status") } /// Enqueue one job and claim it (→ running, `started_at` and `heartbeat_at` /// both stamped NOW()). async fn enqueue_and_claim(h: &TestHarness, user: UserId, key: &str) -> Uuid { scan_jobs::enqueue( &h.db, ScanTargetKind::Item, Uuid::new_v4(), key, FileType::Download, user, 1000, ) .await .expect("enqueue"); let job = scan_jobs::claim_next(&h.db) .await .expect("claim") .expect("a queued job to claim"); assert_eq!(job.status, "running"); assert!(job.heartbeat_at.is_some(), "claim stamps heartbeat_at"); job.id } #[tokio::test] async fn reaper_spares_slow_but_alive_job() { let h = TestHarness::new().await; let user = seed_user(&h, "reaper_alive").await; let job = enqueue_and_claim(&h, user, "k/alive").await; // A job running far past the stuck window but whose worker is still beating: // old started_at, fresh heartbeat. This is the exact case that used to be // reaped and double-processed. sqlx::query( "UPDATE scan_jobs SET started_at = NOW() - interval '1 hour', heartbeat_at = NOW() WHERE id = $1", ) .bind(job) .execute(&h.db) .await .expect("backdate started, fresh beat"); let reaped = scan_jobs::reap_stuck(&h.db, 300).await.expect("reap"); assert_eq!(reaped, 0, "a job with a fresh heartbeat must not be reaped"); assert_eq!(status_of(&h, job).await, "running"); } #[tokio::test] async fn reaper_reaps_crashed_worker() { let h = TestHarness::new().await; let user = seed_user(&h, "reaper_dead").await; let job = enqueue_and_claim(&h, user, "k/dead").await; // Worker crashed: no more beats. sqlx::query("UPDATE scan_jobs SET heartbeat_at = NOW() - interval '1 hour' WHERE id = $1") .bind(job) .execute(&h.db) .await .expect("stale beat"); let reaped = scan_jobs::reap_stuck(&h.db, 300).await.expect("reap"); assert_eq!(reaped, 1, "a stale-heartbeat job must be reaped"); // attempts after a single claim = 1 < MAX_SCAN_ATTEMPTS, so it returns to // the queue for another attempt rather than being retired to failed. assert_eq!(status_of(&h, job).await, "queued"); } #[tokio::test] async fn reaper_falls_back_to_started_at_when_no_heartbeat() { // Any row with heartbeat_at NULL (e.g. one in flight across the migration) // stays reapable on the old started_at clock via COALESCE. let h = TestHarness::new().await; let user = seed_user(&h, "reaper_null").await; let job = enqueue_and_claim(&h, user, "k/null").await; sqlx::query( "UPDATE scan_jobs SET heartbeat_at = NULL, started_at = NOW() - interval '1 hour' WHERE id = $1", ) .bind(job) .execute(&h.db) .await .expect("null beat, old start"); let reaped = scan_jobs::reap_stuck(&h.db, 300).await.expect("reap"); assert_eq!(reaped, 1, "null heartbeat falls back to started_at"); } #[tokio::test] async fn bump_heartbeat_refreshes_and_never_resurrects() { let h = TestHarness::new().await; let user = seed_user(&h, "reaper_bump").await; let job = enqueue_and_claim(&h, user, "k/bump").await; // A stale beat would be reaped... sqlx::query("UPDATE scan_jobs SET heartbeat_at = NOW() - interval '1 hour' WHERE id = $1") .bind(job) .execute(&h.db) .await .expect("stale"); // ...but a bump refreshes it, so the reaper leaves it alone. scan_jobs::bump_heartbeat(&h.db, job).await.expect("bump"); let reaped = scan_jobs::reap_stuck(&h.db, 300).await.expect("reap"); assert_eq!(reaped, 0, "bump_heartbeat must refresh liveness"); assert_eq!(status_of(&h, job).await, "running"); // A late bump after the job leaves `running` is a no-op, it must never // resurrect a terminal row's heartbeat. sqlx::query("UPDATE scan_jobs SET status = 'done', heartbeat_at = NULL WHERE id = $1") .bind(job) .execute(&h.db) .await .expect("mark done"); scan_jobs::bump_heartbeat(&h.db, job) .await .expect("late bump"); let still_null: bool = sqlx::query_scalar("SELECT heartbeat_at IS NULL FROM scan_jobs WHERE id = $1") .bind(job) .fetch_one(&h.db) .await .expect("hb null check"); assert!( still_null, "bump must not touch a job that has left running" ); }