//! DB-layer contract tests for `db/imports.rs`: the import-job CRUD surface, the //! typed status transitions, the liveness heartbeat, and the stuck-job reaper. //! //! Audit A5: imports run on the bounded background pool with no worker loop, so //! before migration 168 a crash mid-import stranded the job in `processing` //! forever. These pin that `update_import_status(Processing)` seeds a heartbeat, //! `bump_import_heartbeat` refreshes it, and `reap_stuck_import_jobs` fails a //! stale-heartbeat job while sparing a freshly-beating or already-terminal one. use crate::harness::TestHarness; use makenotwork::db::{self, ImportJobStatus, ProjectId}; use makenotwork::import::ImportSource; /// Create a creator + project and return (user_id, project_id) typed for direct /// `db::imports` calls. async fn seed_creator(h: &mut TestHarness) -> (makenotwork::db::UserId, ProjectId) { let setup = h .create_creator_with_item("importlayer", "digital", 0) .await; let project_id: ProjectId = setup.project_id.parse().expect("project_id parses"); (setup.user_id, project_id) } #[tokio::test] async fn import_job_crud_roundtrip() { let mut h = TestHarness::new().await; let (user_id, project_id) = seed_creator(&mut h).await; let job = db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 10) .await .expect("create job"); assert_eq!(job.status, ImportJobStatus::Pending); assert_eq!(job.total_rows, 10); db::imports::update_import_progress(&h.db, job.id, 5, 4, 1) .await .expect("progress"); let fetched = db::imports::get_import_job(&h.db, job.id, user_id) .await .expect("get") .expect("job exists"); assert_eq!(fetched.processed_rows, 5); assert_eq!(fetched.created_rows, 4); assert_eq!(fetched.skipped_rows, 1); db::imports::complete_import_job(&h.db, job.id, Some("2 warnings".into())) .await .expect("complete"); let done = db::imports::get_import_job(&h.db, job.id, user_id) .await .expect("get") .expect("job exists"); assert_eq!(done.status, ImportJobStatus::Completed); assert_eq!(done.error_log.as_deref(), Some("2 warnings")); assert!(done.completed_at.is_some()); // Scoping: another user can't read the job. let other = h .signup("importlayer2", "il2@example.com", "Password1!") .await; assert!( db::imports::get_import_job(&h.db, job.id, other) .await .unwrap() .is_none(), "get_import_job must be user-scoped" ); } #[tokio::test] async fn processing_status_seeds_heartbeat() { let mut h = TestHarness::new().await; let (user_id, project_id) = seed_creator(&mut h).await; let job = db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 1) .await .unwrap(); // Pending job has no heartbeat. let hb0: Option> = sqlx::query_scalar("SELECT heartbeat_at FROM import_jobs WHERE id = $1") .bind(job.id) .fetch_one(&h.db) .await .unwrap(); assert!(hb0.is_none(), "a pending job must not have a heartbeat"); db::imports::update_import_status(&h.db, job.id, ImportJobStatus::Processing) .await .unwrap(); let hb1: Option> = sqlx::query_scalar("SELECT heartbeat_at FROM import_jobs WHERE id = $1") .bind(job.id) .fetch_one(&h.db) .await .unwrap(); assert!(hb1.is_some(), "entering Processing must stamp a heartbeat"); } #[tokio::test] async fn reaper_fails_stale_processing_job() { let mut h = TestHarness::new().await; let (user_id, project_id) = seed_creator(&mut h).await; let job = db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 100) .await .unwrap(); db::imports::update_import_status(&h.db, job.id, ImportJobStatus::Processing) .await .unwrap(); // Owning process crashed: heartbeat goes stale. sqlx::query("UPDATE import_jobs SET heartbeat_at = NOW() - interval '1 hour' WHERE id = $1") .bind(job.id) .execute(&h.db) .await .unwrap(); let reaped = db::imports::reap_stuck_import_jobs(&h.db, 1800) .await .unwrap(); assert_eq!(reaped, 1, "a stale-heartbeat processing job must be reaped"); let after = db::imports::get_import_job(&h.db, job.id, user_id) .await .unwrap() .unwrap(); assert_eq!(after.status, ImportJobStatus::Failed); assert!(after.completed_at.is_some()); assert!( after.error_log.as_deref().unwrap_or("").contains("reaped"), "reaped job error_log should explain why: {:?}", after.error_log ); } #[tokio::test] async fn reaper_spares_fresh_and_terminal_jobs() { let mut h = TestHarness::new().await; let (user_id, project_id) = seed_creator(&mut h).await; // Fresh, actively-beating processing job. let fresh = db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 100) .await .unwrap(); db::imports::update_import_status(&h.db, fresh.id, ImportJobStatus::Processing) .await .unwrap(); db::imports::bump_import_heartbeat(&h.db, fresh.id) .await .unwrap(); // Already-completed job, terminal, must never be reaped. let done = db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 1) .await .unwrap(); db::imports::complete_import_job(&h.db, done.id, None) .await .unwrap(); let reaped = db::imports::reap_stuck_import_jobs(&h.db, 1800) .await .unwrap(); assert_eq!( reaped, 0, "a fresh-heartbeat job and a completed job must be spared" ); assert_eq!( db::imports::get_import_job(&h.db, fresh.id, user_id) .await .unwrap() .unwrap() .status, ImportJobStatus::Processing ); assert_eq!( db::imports::get_import_job(&h.db, done.id, user_id) .await .unwrap() .unwrap() .status, ImportJobStatus::Completed ); } #[tokio::test] async fn bump_heartbeat_only_touches_processing_jobs() { let mut h = TestHarness::new().await; let (user_id, project_id) = seed_creator(&mut h).await; // A completed (terminal) job must not be revived to a beating state. let done = db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 1) .await .unwrap(); db::imports::complete_import_job(&h.db, done.id, None) .await .unwrap(); db::imports::bump_import_heartbeat(&h.db, done.id) .await .unwrap(); let hb: Option> = sqlx::query_scalar("SELECT heartbeat_at FROM import_jobs WHERE id = $1") .bind(done.id) .fetch_one(&h.db) .await .unwrap(); assert!( hb.is_none(), "bump must not stamp a heartbeat on a terminal job" ); assert_eq!( db::imports::get_import_job(&h.db, done.id, user_id) .await .unwrap() .unwrap() .status, ImportJobStatus::Completed ); }