Skip to main content

headless_lms_server/programs/
exercise_answer_upload_reaper.rs

1//! Removes files uploaded to be named in an exercise answer that no submission was ever made
2//! from.
3//!
4//! How long an upload is spared depends on its origin, since a native client uploads immediately
5//! before submitting while an iframe student may hold an upload for the length of an exam. The
6//! binding row is soft-deleted rather than removed so that a submit naming a reaped file can
7//! still answer `upload_expired` instead of the misleading `unknown_upload`.
8
9use 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    // The CronJob's exit status is the only signal anyone watches, so a run where every delete
68    // failed must not look green.
69    if failed > 0 {
70        anyhow::bail!(
71            "Failed to reap {failed} of {} answer uploads.",
72            reaped + failed
73        );
74    }
75    Ok(())
76}
77
78/// Retires the binding, removes the object, and only then soft-deletes the `file_uploads` row.
79/// `Ok(false)` means a submission came to reference the upload after `get_reapable` listed it, so
80/// it is no longer reapable.
81///
82/// The order matters in both directions. Retiring the binding first means a submit naming this
83/// upload answers `upload_expired` rather than succeeding and handing the exercise service a URL
84/// that 404s. Deleting the `file_uploads` row last means `get_reapable` still sees the row after a
85/// failed object delete and retries it on a later run, instead of orphaning the object forever.
86async 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    // `optional` tolerates a file a previous interrupted run already soft-deleted, without
96    // swallowing real database errors.
97    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    /// `reap` is global: it takes every eligible row in the shared test database, so two of
113    /// these tests in flight at once reap each other's committed fixtures. Held for the whole test,
114    /// fixtures included. Tokio's mutex rather than `std`'s, so one failing test does not poison
115    /// the lock and fail the rest.
116    static REAPER_TESTS: LazyLock<tokio::sync::Mutex<()>> =
117        LazyLock::new(|| tokio::sync::Mutex::new(()));
118
119    /// Records what the reaper asked it to delete. The real stores are irrelevant here: the
120    /// property under test is which paths the reaper touches.
121    #[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    /// How many times the reaper touched `path`. Scoped this way because the tests commit their
170    /// fixtures, so every run also sees the rows of whatever sibling test is running beside it.
171    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    /// Covers the whole per-row sequence: the object goes, the file row is soft-deleted, and the
181    /// binding survives soft-deleted so submit can still answer `upload_expired`.
182    ///
183    /// Committed, since `reap` now reaps each row over its own pool connection rather than the
184    /// connection fixtures were inserted on.
185    #[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    /// Committed, for the same reason as above.
241    #[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    /// The origin split, end to end: two hours is past a native client's window but nowhere near an
289    /// iframe student's, who may still be holding the file mid-exam.
290    ///
291    /// Committed, for the same reason as the tests above.
292    #[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    /// Committed, for the same reason as the tests above.
339    #[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    /// Fails every delete, standing in for a transient object-store error.
385    #[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    /// A failed object delete must surface in the exit status and be retried later. Before this
438    /// was fixed the run was green and the object was orphaned forever, because the binding's own
439    /// `deleted_at` excluded the row from every later listing.
440    ///
441    /// Committed, for the same reason as the tests above.
442    #[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    /// The reap-vs-submit race with two real connections, which is the part
504    /// `models::exercise_answer_uploads`' own tests cannot reach: they run inside one
505    /// uncommitted transaction, so a second connection can never see their fixtures.
506    ///
507    /// What is under test is not just the outcome but the mechanism — that a reaper running
508    /// concurrently with a submit *blocks* on the row lock `lock_for_exercise_and_user` takes,
509    /// instead of racing past it, and that when it unblocks its `NOT EXISTS` re-check observes the
510    /// association the submit committed. The second half is a Postgres detail worth pinning: the
511    /// blocked `UPDATE` re-evaluates its qual, subquery included, against the committed row, so it
512    /// declines rather than destroying the files of a submission that already returned 200.
513    #[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        // Committed so the reaper's connection can see them. Deliberately not backdated: an
517        // upload inside the retention window is invisible to `get_reapable`, so what this test
518        // leaves in the database cannot perturb the unfiltered `reap()` calls above.
519        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        // The submit side: validate under the row lock, inside the transaction that will record
543        // the association.
544        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        // The reaper, on its own connection, tries to retire the very row the submit holds.
558        let mut reaper_conn = Conn::init().await;
559        let mut reaper_tx = reaper_conn.begin().await;
560        // Scoped so the pinned future releases its borrow of `reaper_tx` before the rollback.
561        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}