Skip to main content

headless_lms_server/programs/
exercise_spec_upload_reaper.rs

1//! Removes files uploaded through the exercise-service upload route that never reached a saved
2//! spec.
3//!
4//! Nothing else can reclaim them: the host stores specs as opaque blobs, so a stored file's only
5//! reference may sit inside content it cannot parse. What makes a file safe here is a declaration
6//! — `exercise_task_spec_files` for a live spec, `page_history_spec_files` for every snapshot a
7//! restore could bring back. Because history is kept, this reaper mostly collects uploads a
8//! teacher abandoned before saving, not files dropped from a spec that was once saved.
9//!
10//! The binding row is soft-deleted rather than removed, leaving an audit trail of what was
11//! reclaimed.
12
13use std::path::Path;
14
15use crate::config::{FileStoreRuntimeConfig, program_config::ProgramConfig};
16use crate::{setup_file_store, setup_tracing};
17use dotenvy::dotenv;
18use futures::{StreamExt, stream};
19use headless_lms_models::{self as models, error::TryToOptional};
20use headless_lms_utils::file_store::FileStore;
21use sqlx::{PgConnection, PgPool};
22
23const MAX_CONCURRENT_REAPS: usize = 8;
24
25pub async fn main() -> anyhow::Result<()> {
26    dotenv().ok();
27    ProgramConfig::ensure_default_rust_log_for_workers();
28    setup_tracing()?;
29    let database_url = ProgramConfig::database_url_with_default();
30    let base_url = ProgramConfig::required("BASE_URL")?;
31    let file_store = setup_file_store(&FileStoreRuntimeConfig::try_from_env()?, &base_url).await;
32    let db_pool = PgPool::connect(&database_url).await?;
33    reap(&db_pool, file_store.as_ref()).await
34}
35
36async fn reap(pool: &PgPool, file_store: &dyn FileStore) -> anyhow::Result<()> {
37    let mut conn = pool.acquire().await?;
38    let reapable = models::exercise_spec_uploads::get_reapable(&mut conn).await?;
39    drop(conn);
40    info!("Reaping {} abandoned spec uploads.", reapable.len());
41
42    let mut reaped = 0;
43    let mut skipped = 0;
44    let mut failed = 0;
45    let mut results = stream::iter(reapable)
46        .map(|upload| async move {
47            let file_upload_id = upload.file_upload_id;
48            let result = match pool.acquire().await {
49                Ok(mut conn) => reap_one(&mut conn, file_store, &upload).await,
50                Err(err) => Err(err.into()),
51            };
52            (file_upload_id, result)
53        })
54        .buffer_unordered(MAX_CONCURRENT_REAPS);
55    while let Some((file_upload_id, result)) = results.next().await {
56        match result {
57            Ok(true) => reaped += 1,
58            Ok(false) => skipped += 1,
59            Err(err) => {
60                failed += 1;
61                error!("Failed to reap spec upload {}: {:#?}", file_upload_id, err);
62            }
63        }
64    }
65    info!(
66        "Abandoned spec uploads reaped. Succeeded: {reaped}, skipped: {skipped}, failed: {failed}."
67    );
68    // The CronJob's exit status is the only signal anyone watches, so a run where every delete
69    // failed must not look green.
70    if failed > 0 {
71        anyhow::bail!(
72            "Failed to reap {failed} of {} spec uploads.",
73            reaped + failed
74        );
75    }
76    Ok(())
77}
78
79/// Retires the record, removes the object, and only then soft-deletes the `file_uploads` row.
80/// `Ok(false)` means a save came to declare the file after `get_reapable` listed it, so it is no
81/// longer reapable.
82///
83/// Deleting the `file_uploads` row last is what makes a failed object delete recoverable:
84/// `get_reapable` still sees the row and retries it on a later run, instead of orphaning the
85/// object forever.
86async fn reap_one(
87    conn: &mut PgConnection,
88    file_store: &dyn FileStore,
89    upload: &models::exercise_spec_uploads::ReapableUpload,
90) -> anyhow::Result<bool> {
91    if !models::exercise_spec_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}