1use std::path::Path;
10
11use crate::config::{FileStoreRuntimeConfig, program_config::ProgramConfig};
12use crate::{setup_file_store, setup_tracing};
13use dotenvy::dotenv;
14use futures::{StreamExt, stream};
15use headless_lms_models::{self as models, error::TryToOptional};
16use headless_lms_utils::file_store::FileStore;
17use sqlx::{PgConnection, PgPool};
18
19const MAX_CONCURRENT_REAPS: usize = 8;
20
21pub async fn main() -> anyhow::Result<()> {
22 dotenv().ok();
23 ProgramConfig::ensure_default_rust_log_for_workers();
24 setup_tracing()?;
25 let database_url = ProgramConfig::database_url_with_default();
26 let base_url = ProgramConfig::required("BASE_URL")?;
27 let file_store = setup_file_store(&FileStoreRuntimeConfig::try_from_env()?, &base_url).await;
28 let db_pool = PgPool::connect(&database_url).await?;
29 reap(&db_pool, file_store.as_ref()).await
30}
31
32async fn reap(pool: &PgPool, file_store: &dyn FileStore) -> anyhow::Result<()> {
33 let mut conn = pool.acquire().await?;
34 let reapable = models::exercise_answer_uploads::get_reapable(&mut conn).await?;
35 drop(conn);
36 info!("Reaping {} orphaned answer uploads.", reapable.len());
37
38 let mut reaped = 0;
39 let mut skipped = 0;
40 let mut failed = 0;
41 let mut results = stream::iter(reapable)
42 .map(|upload| async move {
43 let file_upload_id = upload.file_upload_id;
44 let result = match pool.acquire().await {
45 Ok(mut conn) => reap_one(&mut conn, file_store, &upload).await,
46 Err(err) => Err(err.into()),
47 };
48 (file_upload_id, result)
49 })
50 .buffer_unordered(MAX_CONCURRENT_REAPS);
51 while let Some((file_upload_id, result)) = results.next().await {
52 match result {
53 Ok(true) => reaped += 1,
54 Ok(false) => skipped += 1,
55 Err(err) => {
56 failed += 1;
57 error!(
58 "Failed to reap answer upload {}: {:#?}",
59 file_upload_id, err
60 );
61 }
62 }
63 }
64 info!(
65 "Orphaned answer uploads reaped. Succeeded: {reaped}, skipped: {skipped}, failed: {failed}."
66 );
67 if failed > 0 {
70 anyhow::bail!(
71 "Failed to reap {failed} of {} answer uploads.",
72 reaped + failed
73 );
74 }
75 Ok(())
76}
77
78async fn reap_one(
87 conn: &mut PgConnection,
88 file_store: &dyn FileStore,
89 upload: &models::exercise_answer_uploads::ReapableUpload,
90) -> anyhow::Result<bool> {
91 if !models::exercise_answer_uploads::mark_reaped(conn, upload.id).await? {
92 return Ok(false);
93 }
94 file_store.delete(Path::new(&upload.path)).await?;
95 models::file_uploads::delete_and_fetch_path(conn, upload.file_upload_id)
98 .await
99 .optional()?;
100 Ok(true)
101}
102
103#[cfg(test)]
104mod tests {
105 use super::*;
106 use crate::test_helper::*;
107 use chrono::Duration;
108 use headless_lms_base::error::backend_error::BackendError;
109 use headless_lms_utils::prelude::{UtilError, UtilErrorType, UtilResult};
110 use std::sync::{LazyLock, Mutex};
111
112 static REAPER_TESTS: LazyLock<tokio::sync::Mutex<()>> =
117 LazyLock::new(|| tokio::sync::Mutex::new(()));
118
119 #[derive(Default)]
122 struct RecordingFileStore {
123 deleted: Mutex<Vec<String>>,
124 }
125
126 #[async_trait::async_trait(?Send)]
127 impl FileStore for RecordingFileStore {
128 async fn upload(&self, _path: &Path, _contents: Vec<u8>, _mime: &str) -> UtilResult<()> {
129 unimplemented!("not reached by the reaper")
130 }
131
132 async fn upload_stream(
133 &self,
134 _path: &Path,
135 _contents: headless_lms_utils::file_store::GenericPayload,
136 _mime: &str,
137 ) -> UtilResult<()> {
138 unimplemented!("not reached by the reaper")
139 }
140
141 async fn download(&self, _path: &Path) -> UtilResult<Vec<u8>> {
142 unimplemented!("not reached by the reaper")
143 }
144
145 async fn download_stream(
146 &self,
147 _path: &Path,
148 ) -> UtilResult<Box<dyn futures::Stream<Item = std::io::Result<bytes::Bytes>>>> {
149 unimplemented!("not reached by the reaper")
150 }
151
152 async fn get_direct_download_url(&self, _path: &Path) -> UtilResult<String> {
153 unimplemented!("not reached by the reaper")
154 }
155
156 async fn delete(&self, path: &Path) -> UtilResult<()> {
157 self.deleted
158 .lock()
159 .expect("lock")
160 .push(path.to_string_lossy().to_string());
161 Ok(())
162 }
163
164 fn get_cache_files_folder_path(&self) -> UtilResult<&Path> {
165 unimplemented!("not reached by the reaper")
166 }
167 }
168
169 fn deletions_of(recorded: &Mutex<Vec<String>>, path: &str) -> usize {
172 recorded
173 .lock()
174 .expect("lock")
175 .iter()
176 .filter(|recorded_path| recorded_path.as_str() == path)
177 .count()
178 }
179
180 #[actix_web::test]
186 async fn reaps_an_orphaned_upload_and_leaves_a_row_behind() {
187 let _serialized = REAPER_TESTS.lock().await;
188 insert_data!(:tx, user: user, :org, :course, instance: _instance, :course_module, :chapter, :page, :exercise, :slide, task: _task);
189 let path = "exercise-services-client/orphan";
190 let file_id = models::file_uploads::insert(
191 tx.as_mut(),
192 "orphan.tar.zst",
193 path,
194 "application/octet-stream",
195 Some(user),
196 None,
197 )
198 .await
199 .expect("file upload");
200 models::exercise_answer_uploads::insert_many(
201 tx.as_mut(),
202 exercise,
203 user,
204 &[file_id],
205 models::exercise_answer_uploads::AnswerUploadOrigin::NativeClient,
206 )
207 .await
208 .expect("binding");
209 backdate(tx.as_mut(), file_id, Duration::hours(2)).await;
210 tx.commit().await;
211
212 let pool = PgPool::connect(&test_database_url())
213 .await
214 .expect("test pool");
215 let file_store = RecordingFileStore::default();
216 reap(&pool, &file_store).await.expect("reap");
217
218 assert_eq!(deletions_of(&file_store.deleted, path), 1);
219 let mut check_conn = Conn::init().await;
220 let mut check_tx = check_conn.begin().await;
221 assert!(
222 models::file_uploads::get_many(check_tx.as_mut(), &[file_id])
223 .await
224 .expect("file lookup")
225 .is_empty()
226 );
227 let recorded = models::exercise_answer_uploads::get_for_exercise_and_user(
228 check_tx.as_mut(),
229 exercise,
230 user,
231 &[file_id],
232 )
233 .await
234 .expect("binding lookup");
235 assert_eq!(recorded.len(), 1);
236 assert!(recorded[0].deleted);
237 check_tx.rollback().await;
238 }
239
240 #[actix_web::test]
242 async fn spares_an_upload_inside_the_retention_window() {
243 let _serialized = REAPER_TESTS.lock().await;
244 insert_data!(:tx, user: user, :org, :course, instance: _instance, :course_module, :chapter, :page, :exercise, :slide, task: _task);
245 let file_id = models::file_uploads::insert(
246 tx.as_mut(),
247 "fresh.tar.zst",
248 "exercise-services-client/fresh",
249 "application/octet-stream",
250 Some(user),
251 None,
252 )
253 .await
254 .expect("file upload");
255 models::exercise_answer_uploads::insert_many(
256 tx.as_mut(),
257 exercise,
258 user,
259 &[file_id],
260 models::exercise_answer_uploads::AnswerUploadOrigin::NativeClient,
261 )
262 .await
263 .expect("binding");
264 tx.commit().await;
265
266 let pool = PgPool::connect(&test_database_url())
267 .await
268 .expect("test pool");
269 let file_store = RecordingFileStore::default();
270 reap(&pool, &file_store).await.expect("reap");
271
272 assert_eq!(
273 deletions_of(&file_store.deleted, "exercise-services-client/fresh"),
274 0
275 );
276 let mut check_conn = Conn::init().await;
277 let mut check_tx = check_conn.begin().await;
278 assert_eq!(
279 models::file_uploads::get_many(check_tx.as_mut(), &[file_id])
280 .await
281 .expect("file lookup")
282 .len(),
283 1
284 );
285 check_tx.rollback().await;
286 }
287
288 #[actix_web::test]
293 async fn spares_an_iframe_upload_a_native_client_upload_would_lose() {
294 let _serialized = REAPER_TESTS.lock().await;
295 insert_data!(:tx, user: user, :org, :course, instance: _instance, :course_module, :chapter, :page, :exercise, :slide, task: _task);
296 let path = "exercise-services-client/mid-exam";
297 let file_id = models::file_uploads::insert(
298 tx.as_mut(),
299 "mid-exam.pdf",
300 path,
301 "application/pdf",
302 Some(user),
303 None,
304 )
305 .await
306 .expect("file upload");
307 models::exercise_answer_uploads::insert_many(
308 tx.as_mut(),
309 exercise,
310 user,
311 &[file_id],
312 models::exercise_answer_uploads::AnswerUploadOrigin::Iframe,
313 )
314 .await
315 .expect("binding");
316 backdate(tx.as_mut(), file_id, Duration::hours(2)).await;
317 tx.commit().await;
318
319 let pool = PgPool::connect(&test_database_url())
320 .await
321 .expect("test pool");
322 let file_store = RecordingFileStore::default();
323 reap(&pool, &file_store).await.expect("reap");
324
325 assert_eq!(deletions_of(&file_store.deleted, path), 0);
326 let mut check_conn = Conn::init().await;
327 let mut check_tx = check_conn.begin().await;
328 assert_eq!(
329 models::file_uploads::get_many(check_tx.as_mut(), &[file_id])
330 .await
331 .expect("file lookup")
332 .len(),
333 1
334 );
335 check_tx.rollback().await;
336 }
337
338 #[actix_web::test]
340 async fn reaps_an_iframe_upload_past_seven_days() {
341 let _serialized = REAPER_TESTS.lock().await;
342 insert_data!(:tx, user: user, :org, :course, instance: _instance, :course_module, :chapter, :page, :exercise, :slide, task: _task);
343 let path = "exercise-services-client/abandoned";
344 let file_id = models::file_uploads::insert(
345 tx.as_mut(),
346 "abandoned.pdf",
347 path,
348 "application/pdf",
349 Some(user),
350 None,
351 )
352 .await
353 .expect("file upload");
354 models::exercise_answer_uploads::insert_many(
355 tx.as_mut(),
356 exercise,
357 user,
358 &[file_id],
359 models::exercise_answer_uploads::AnswerUploadOrigin::Iframe,
360 )
361 .await
362 .expect("binding");
363 backdate(tx.as_mut(), file_id, Duration::days(8)).await;
364 tx.commit().await;
365
366 let pool = PgPool::connect(&test_database_url())
367 .await
368 .expect("test pool");
369 let file_store = RecordingFileStore::default();
370 reap(&pool, &file_store).await.expect("reap");
371
372 assert_eq!(deletions_of(&file_store.deleted, path), 1);
373 let mut check_conn = Conn::init().await;
374 let mut check_tx = check_conn.begin().await;
375 assert!(
376 models::file_uploads::get_many(check_tx.as_mut(), &[file_id])
377 .await
378 .expect("file lookup")
379 .is_empty()
380 );
381 check_tx.rollback().await;
382 }
383
384 #[derive(Default)]
386 struct FailingFileStore {
387 attempts: Mutex<Vec<String>>,
388 }
389
390 #[async_trait::async_trait(?Send)]
391 impl FileStore for FailingFileStore {
392 async fn upload(&self, _path: &Path, _contents: Vec<u8>, _mime: &str) -> UtilResult<()> {
393 unimplemented!("not reached by the reaper")
394 }
395
396 async fn upload_stream(
397 &self,
398 _path: &Path,
399 _contents: headless_lms_utils::file_store::GenericPayload,
400 _mime: &str,
401 ) -> UtilResult<()> {
402 unimplemented!("not reached by the reaper")
403 }
404
405 async fn download(&self, _path: &Path) -> UtilResult<Vec<u8>> {
406 unimplemented!("not reached by the reaper")
407 }
408
409 async fn download_stream(
410 &self,
411 _path: &Path,
412 ) -> UtilResult<Box<dyn futures::Stream<Item = std::io::Result<bytes::Bytes>>>> {
413 unimplemented!("not reached by the reaper")
414 }
415
416 async fn get_direct_download_url(&self, _path: &Path) -> UtilResult<String> {
417 unimplemented!("not reached by the reaper")
418 }
419
420 async fn delete(&self, path: &Path) -> UtilResult<()> {
421 self.attempts
422 .lock()
423 .expect("lock")
424 .push(path.to_string_lossy().to_string());
425 Err(UtilError::new(
426 UtilErrorType::Other,
427 "simulated object store failure".to_string(),
428 None,
429 ))
430 }
431
432 fn get_cache_files_folder_path(&self) -> UtilResult<&Path> {
433 unimplemented!("not reached by the reaper")
434 }
435 }
436
437 #[actix_web::test]
443 async fn a_failed_object_delete_fails_the_run_and_is_retried() {
444 let _serialized = REAPER_TESTS.lock().await;
445 insert_data!(:tx, user: user, :org, :course, instance: _instance, :course_module, :chapter, :page, :exercise, :slide, task: _task);
446 let path = "exercise-services-client/transient";
447 let file_id = models::file_uploads::insert(
448 tx.as_mut(),
449 "transient.tar.zst",
450 path,
451 "application/octet-stream",
452 Some(user),
453 None,
454 )
455 .await
456 .expect("file upload");
457 models::exercise_answer_uploads::insert_many(
458 tx.as_mut(),
459 exercise,
460 user,
461 &[file_id],
462 models::exercise_answer_uploads::AnswerUploadOrigin::NativeClient,
463 )
464 .await
465 .expect("binding");
466 backdate(tx.as_mut(), file_id, Duration::hours(2)).await;
467 tx.commit().await;
468
469 let pool = PgPool::connect(&test_database_url())
470 .await
471 .expect("test pool");
472 let failing = FailingFileStore::default();
473 reap(&pool, &failing)
474 .await
475 .expect_err("a run where every delete failed must not exit successfully");
476 assert_eq!(deletions_of(&failing.attempts, path), 1);
477 let mut check_conn = Conn::init().await;
478 let mut check_tx = check_conn.begin().await;
479 assert_eq!(
480 models::file_uploads::get_many(check_tx.as_mut(), &[file_id])
481 .await
482 .expect("file lookup")
483 .len(),
484 1,
485 "the file row must survive, since it is what makes the retry possible"
486 );
487 check_tx.rollback().await;
488
489 let retry = RecordingFileStore::default();
490 reap(&pool, &retry).await.expect("retry");
491 assert_eq!(deletions_of(&retry.deleted, path), 1);
492 let mut check_conn = Conn::init().await;
493 let mut check_tx = check_conn.begin().await;
494 assert!(
495 models::file_uploads::get_many(check_tx.as_mut(), &[file_id])
496 .await
497 .expect("file lookup")
498 .is_empty()
499 );
500 check_tx.rollback().await;
501 }
502
503 #[actix_web::test]
514 async fn a_concurrent_reaper_blocks_on_the_submit_lock_and_then_declines_to_reap() {
515 let _serialized = REAPER_TESTS.lock().await;
516 insert_data!(:tx, user: user, :org, course: course, instance: _instance, :course_module, :chapter, :page, :exercise, slide: slide, task: task);
520 let file_id = models::file_uploads::insert(
521 tx.as_mut(),
522 "raced.tar.zst",
523 "exercise-services-client/raced",
524 "application/octet-stream",
525 Some(user),
526 None,
527 )
528 .await
529 .expect("file upload");
530 models::exercise_answer_uploads::insert_many(
531 tx.as_mut(),
532 exercise,
533 user,
534 &[file_id],
535 models::exercise_answer_uploads::AnswerUploadOrigin::NativeClient,
536 )
537 .await
538 .expect("binding");
539 let binding_id = binding_id_of(tx.as_mut(), file_id).await;
540 tx.commit().await;
541
542 let mut submit_conn = Conn::init().await;
545 let mut submit_tx = submit_conn.begin().await;
546 let locked = models::exercise_answer_uploads::lock_for_exercise_and_user(
547 submit_tx.as_mut(),
548 exercise,
549 user,
550 &[file_id],
551 )
552 .await
553 .expect("locked lookup");
554 assert_eq!(locked.len(), 1);
555 assert!(!locked[0].deleted);
556
557 let mut reaper_conn = Conn::init().await;
559 let mut reaper_tx = reaper_conn.begin().await;
560 let reaped = {
562 let mut reap = std::pin::pin!(models::exercise_answer_uploads::mark_reaped(
563 reaper_tx.as_mut(),
564 binding_id
565 ));
566 assert!(
567 tokio::time::timeout(std::time::Duration::from_millis(500), &mut reap)
568 .await
569 .is_err(),
570 "the reaper must block on the row lock the submit holds, not decide without it"
571 );
572
573 let slide_submission = models::exercise_slide_submissions::insert_exercise_slide_submission(
574 submit_tx.as_mut(),
575 models::exercise_slide_submissions::NewExerciseSlideSubmission {
576 exercise_slide_id: slide,
577 course_id: Some(course),
578 exam_id: None,
579 user_id: user,
580 exercise_id: exercise,
581 user_points_update_strategy:
582 models::exercise_task_gradings::UserPointsUpdateStrategy::CanAddPointsAndCanRemovePoints,
583 },
584 )
585 .await
586 .expect("slide submission");
587 let task_submission = models::exercise_task_submissions::insert(
588 submit_tx.as_mut(),
589 models::PKeyPolicy::Generate,
590 slide_submission.id,
591 slide,
592 task,
593 &models::library::grading::SubmittedAnswer::Json {
594 data: serde_json::json!({ "opaque": "plugin owned" }),
595 },
596 )
597 .await
598 .expect("task submission");
599 models::exercise_task_submission_files::insert_many(
600 submit_tx.as_mut(),
601 task_submission,
602 &[file_id],
603 )
604 .await
605 .expect("submission files");
606 submit_tx.commit().await;
607
608 tokio::time::timeout(std::time::Duration::from_secs(10), &mut reap)
609 .await
610 .expect("the reaper must unblock once the submit commits")
611 .expect("mark_reaped")
612 };
613 assert!(
614 !reaped,
615 "the reaper must decline an upload the submit referenced while it waited"
616 );
617 reaper_tx.rollback().await;
618
619 let mut check_conn = Conn::init().await;
620 let mut check_tx = check_conn.begin().await;
621 let recorded = models::exercise_answer_uploads::get_for_exercise_and_user(
622 check_tx.as_mut(),
623 exercise,
624 user,
625 &[file_id],
626 )
627 .await
628 .expect("binding lookup");
629 assert_eq!(
630 recorded,
631 vec![models::exercise_answer_uploads::AnswerUpload {
632 file_upload_id: file_id,
633 deleted: false
634 }],
635 "the upload must stay usable, so download_submission can still serve it"
636 );
637 check_tx.rollback().await;
638 }
639
640 async fn binding_id_of(conn: &mut PgConnection, file_upload_id: uuid::Uuid) -> uuid::Uuid {
641 models::exercise_answer_uploads::get_id_by_file_upload_id(conn, file_upload_id)
642 .await
643 .expect("binding id")
644 }
645
646 async fn backdate(conn: &mut PgConnection, file_upload_id: uuid::Uuid, age: Duration) {
647 models::exercise_answer_uploads::backdate(conn, file_upload_id, age)
648 .await
649 .expect("backdate");
650 }
651}