Skip to main content

max / makenotwork

7.4 KB · 226 lines History Blame Raw
1 //! DB-layer contract tests for `db/imports.rs`: the import-job CRUD surface, the
2 //! typed status transitions, the liveness heartbeat, and the stuck-job reaper.
3 //!
4 //! Audit A5: imports run on the bounded background pool with no worker loop, so
5 //! before migration 168 a crash mid-import stranded the job in `processing`
6 //! forever. These pin that `update_import_status(Processing)` seeds a heartbeat,
7 //! `bump_import_heartbeat` refreshes it, and `reap_stuck_import_jobs` fails a
8 //! stale-heartbeat job while sparing a freshly-beating or already-terminal one.
9
10 use crate::harness::TestHarness;
11 use makenotwork::db::{self, ImportJobStatus, ProjectId};
12 use makenotwork::import::ImportSource;
13
14 /// Create a creator + project and return (user_id, project_id) typed for direct
15 /// `db::imports` calls.
16 async fn seed_creator(h: &mut TestHarness) -> (makenotwork::db::UserId, ProjectId) {
17 let setup = h
18 .create_creator_with_item("importlayer", "digital", 0)
19 .await;
20 let project_id: ProjectId = setup.project_id.parse().expect("project_id parses");
21 (setup.user_id, project_id)
22 }
23
24 #[tokio::test]
25 async fn import_job_crud_roundtrip() {
26 let mut h = TestHarness::new().await;
27 let (user_id, project_id) = seed_creator(&mut h).await;
28
29 let job =
30 db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 10)
31 .await
32 .expect("create job");
33 assert_eq!(job.status, ImportJobStatus::Pending);
34 assert_eq!(job.total_rows, 10);
35
36 db::imports::update_import_progress(&h.db, job.id, 5, 4, 1)
37 .await
38 .expect("progress");
39 let fetched = db::imports::get_import_job(&h.db, job.id, user_id)
40 .await
41 .expect("get")
42 .expect("job exists");
43 assert_eq!(fetched.processed_rows, 5);
44 assert_eq!(fetched.created_rows, 4);
45 assert_eq!(fetched.skipped_rows, 1);
46
47 db::imports::complete_import_job(&h.db, job.id, Some("2 warnings".into()))
48 .await
49 .expect("complete");
50 let done = db::imports::get_import_job(&h.db, job.id, user_id)
51 .await
52 .expect("get")
53 .expect("job exists");
54 assert_eq!(done.status, ImportJobStatus::Completed);
55 assert_eq!(done.error_log.as_deref(), Some("2 warnings"));
56 assert!(done.completed_at.is_some());
57
58 // Scoping: another user can't read the job.
59 let other = h
60 .signup("importlayer2", "il2@example.com", "Password1!")
61 .await;
62 assert!(
63 db::imports::get_import_job(&h.db, job.id, other)
64 .await
65 .unwrap()
66 .is_none(),
67 "get_import_job must be user-scoped"
68 );
69 }
70
71 #[tokio::test]
72 async fn processing_status_seeds_heartbeat() {
73 let mut h = TestHarness::new().await;
74 let (user_id, project_id) = seed_creator(&mut h).await;
75 let job =
76 db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 1)
77 .await
78 .unwrap();
79
80 // Pending job has no heartbeat.
81 let hb0: Option<chrono::DateTime<chrono::Utc>> =
82 sqlx::query_scalar("SELECT heartbeat_at FROM import_jobs WHERE id = $1")
83 .bind(job.id)
84 .fetch_one(&h.db)
85 .await
86 .unwrap();
87 assert!(hb0.is_none(), "a pending job must not have a heartbeat");
88
89 db::imports::update_import_status(&h.db, job.id, ImportJobStatus::Processing)
90 .await
91 .unwrap();
92 let hb1: Option<chrono::DateTime<chrono::Utc>> =
93 sqlx::query_scalar("SELECT heartbeat_at FROM import_jobs WHERE id = $1")
94 .bind(job.id)
95 .fetch_one(&h.db)
96 .await
97 .unwrap();
98 assert!(hb1.is_some(), "entering Processing must stamp a heartbeat");
99 }
100
101 #[tokio::test]
102 async fn reaper_fails_stale_processing_job() {
103 let mut h = TestHarness::new().await;
104 let (user_id, project_id) = seed_creator(&mut h).await;
105 let job =
106 db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 100)
107 .await
108 .unwrap();
109 db::imports::update_import_status(&h.db, job.id, ImportJobStatus::Processing)
110 .await
111 .unwrap();
112
113 // Owning process crashed: heartbeat goes stale.
114 sqlx::query("UPDATE import_jobs SET heartbeat_at = NOW() - interval '1 hour' WHERE id = $1")
115 .bind(job.id)
116 .execute(&h.db)
117 .await
118 .unwrap();
119
120 let reaped = db::imports::reap_stuck_import_jobs(&h.db, 1800)
121 .await
122 .unwrap();
123 assert_eq!(reaped, 1, "a stale-heartbeat processing job must be reaped");
124
125 let after = db::imports::get_import_job(&h.db, job.id, user_id)
126 .await
127 .unwrap()
128 .unwrap();
129 assert_eq!(after.status, ImportJobStatus::Failed);
130 assert!(after.completed_at.is_some());
131 assert!(
132 after.error_log.as_deref().unwrap_or("").contains("reaped"),
133 "reaped job error_log should explain why: {:?}",
134 after.error_log
135 );
136 }
137
138 #[tokio::test]
139 async fn reaper_spares_fresh_and_terminal_jobs() {
140 let mut h = TestHarness::new().await;
141 let (user_id, project_id) = seed_creator(&mut h).await;
142
143 // Fresh, actively-beating processing job.
144 let fresh =
145 db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 100)
146 .await
147 .unwrap();
148 db::imports::update_import_status(&h.db, fresh.id, ImportJobStatus::Processing)
149 .await
150 .unwrap();
151 db::imports::bump_import_heartbeat(&h.db, fresh.id)
152 .await
153 .unwrap();
154
155 // Already-completed job, terminal, must never be reaped.
156 let done =
157 db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 1)
158 .await
159 .unwrap();
160 db::imports::complete_import_job(&h.db, done.id, None)
161 .await
162 .unwrap();
163
164 let reaped = db::imports::reap_stuck_import_jobs(&h.db, 1800)
165 .await
166 .unwrap();
167 assert_eq!(
168 reaped, 0,
169 "a fresh-heartbeat job and a completed job must be spared"
170 );
171
172 assert_eq!(
173 db::imports::get_import_job(&h.db, fresh.id, user_id)
174 .await
175 .unwrap()
176 .unwrap()
177 .status,
178 ImportJobStatus::Processing
179 );
180 assert_eq!(
181 db::imports::get_import_job(&h.db, done.id, user_id)
182 .await
183 .unwrap()
184 .unwrap()
185 .status,
186 ImportJobStatus::Completed
187 );
188 }
189
190 #[tokio::test]
191 async fn bump_heartbeat_only_touches_processing_jobs() {
192 let mut h = TestHarness::new().await;
193 let (user_id, project_id) = seed_creator(&mut h).await;
194
195 // A completed (terminal) job must not be revived to a beating state.
196 let done =
197 db::imports::create_import_job(&h.db, user_id, project_id, ImportSource::GenericCsv, 1)
198 .await
199 .unwrap();
200 db::imports::complete_import_job(&h.db, done.id, None)
201 .await
202 .unwrap();
203 db::imports::bump_import_heartbeat(&h.db, done.id)
204 .await
205 .unwrap();
206
207 let hb: Option<chrono::DateTime<chrono::Utc>> =
208 sqlx::query_scalar("SELECT heartbeat_at FROM import_jobs WHERE id = $1")
209 .bind(done.id)
210 .fetch_one(&h.db)
211 .await
212 .unwrap();
213 assert!(
214 hb.is_none(),
215 "bump must not stamp a heartbeat on a terminal job"
216 );
217 assert_eq!(
218 db::imports::get_import_job(&h.db, done.id, user_id)
219 .await
220 .unwrap()
221 .unwrap()
222 .status,
223 ImportJobStatus::Completed
224 );
225 }
226