Skip to main content

headless_lms_server/programs/
ended_exams_processor.rs

1use std::collections::{HashMap, HashSet};
2
3use crate::config::program_config::ProgramConfig;
4use crate::setup_tracing;
5use chrono::{Duration, Utc};
6use dotenvy::dotenv;
7use headless_lms_base::error::backend_error::BackendError;
8use headless_lms_models::{self as models, ModelError, ModelErrorType};
9use sqlx::{Connection, PgConnection, PgPool};
10use uuid::Uuid;
11
12pub async fn main() -> anyhow::Result<()> {
13    dotenv().ok();
14    ProgramConfig::ensure_default_rust_log_for_workers();
15    setup_tracing()?;
16    let database_url = ProgramConfig::database_url_with_default();
17    let db_pool = PgPool::connect(&database_url).await?;
18    let mut conn = db_pool.acquire().await?;
19    process_ended_exams(&mut conn).await?;
20    process_ended_exam_enrollments(&mut conn).await
21}
22
23/// Fetches ended exams that haven't yet been processed and updates completions for them.
24async fn process_ended_exams(conn: &mut sqlx::PgConnection) -> anyhow::Result<()> {
25    let now = Utc::now();
26    let exam_ids =
27        models::ended_processed_exams::get_unprocessed_ended_exams_by_timestamp(conn, now).await?;
28    tracing::info!("Processing completions for {} ended exams.", exam_ids.len());
29    let mut processed_courses_cache = HashSet::new();
30    let mut success = 0;
31    for exam_id in exam_ids.iter() {
32        match process_ended_exam(conn, *exam_id, &mut processed_courses_cache).await {
33            Ok(_) => success += 1,
34            Err(err) => {
35                tracing::error!("Failed to process exam {}: {:#?}", exam_id, err);
36            }
37        }
38    }
39    tracing::info!(
40        "Exams processed. Succeeded: {}, failed: {}.",
41        success,
42        exam_ids.len() - success
43    );
44    Ok(())
45}
46
47/// Processes completions for courses associated with the given exam.
48///
49/// Because the same course can belong to multiple exams at the same time, a cache for already
50/// processed courses can be provided to avoid unnecessarily reprocessing those courses again.
51async fn process_ended_exam(
52    conn: &mut PgConnection,
53    exam_id: Uuid,
54    already_processed_courses: &mut HashSet<Uuid>,
55) -> anyhow::Result<()> {
56    let course_ids = models::course_exams::get_course_ids_by_exam_id(conn, exam_id).await?;
57    let mut tx = conn.begin().await?;
58    for course_id in course_ids {
59        if already_processed_courses.contains(&course_id) {
60            continue;
61        } else {
62            models::library::progressing::process_all_course_completions(&mut tx, course_id)
63                .await?;
64            already_processed_courses.insert(course_id);
65        }
66    }
67    models::ended_processed_exams::upsert(&mut tx, exam_id).await?;
68    tx.commit().await?;
69    Ok(())
70}
71
72/// Processes ended exam enrollments
73async fn process_ended_exam_enrollments(conn: &mut PgConnection) -> anyhow::Result<()> {
74    let mut tx = conn.begin().await?;
75    let mut success = 0;
76    let mut failed = 0;
77
78    let ongoing_exam_enrollments: Vec<headless_lms_models::exams::ExamEnrollment> =
79        models::exams::get_ongoing_exam_enrollments(&mut tx).await?;
80    let exams = models::exams::get_exams(&mut tx).await?;
81
82    let mut needs_ended_at_date: HashMap<Uuid, Vec<Uuid>> = HashMap::new();
83    for enrollment in ongoing_exam_enrollments {
84        let exam = exams
85            .get(&enrollment.exam_id)
86            .ok_or_else(|| ModelError::new(ModelErrorType::Generic, "Exam not found", None))?;
87
88        //Check if users exams should have ended
89        if Utc::now() > enrollment.started_at + Duration::minutes(exam.time_minutes.into()) {
90            needs_ended_at_date
91                .entry(exam.id)
92                .or_default()
93                .push(enrollment.user_id);
94        }
95    }
96
97    for entry in needs_ended_at_date.into_iter() {
98        let exam_id = entry.0;
99        let user_ids = entry.1;
100        match models::exams::update_exam_ended_at_for_users_with_exam_id(
101            &mut tx,
102            exam_id,
103            &user_ids,
104            Utc::now(),
105        )
106        .await
107        {
108            Ok(_) => success += user_ids.len(),
109            Err(err) => {
110                failed += user_ids.len();
111                tracing::error!(
112                    "Failed to end exam enrolments for exam {}: {:#?}",
113                    exam_id,
114                    err
115                );
116            }
117        }
118    }
119
120    tracing::info!(
121        "Exam enrollments processed. Succeeded: {}, failed: {}.",
122        success,
123        failed
124    );
125    tx.commit().await?;
126
127    Ok(())
128}