headless_lms_server/programs/
exercise_spec_upload_reaper.rs1use 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 if failed > 0 {
71 anyhow::bail!(
72 "Failed to reap {failed} of {} spec uploads.",
73 reaped + failed
74 );
75 }
76 Ok(())
77}
78
79async 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 models::file_uploads::delete_and_fetch_path(conn, upload.file_upload_id)
98 .await
99 .optional()?;
100 Ok(true)
101}