headless_lms_server/programs/
ended_exams_processor.rs1use 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
23async 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
47async 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
72async 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 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}