1use chrono::Duration;
9
10use crate::prelude::*;
11use headless_lms_credit_registration::{
12 CreditRegistrationPhase, PhaseContext, PhaseSkipReason, PhaseTick, Runner, registry_health,
13 run_phase_once,
14};
15use headless_lms_models::credit_registrations::RegistrationScope;
16use headless_lms_models::library::credit_registration::enrolment_check_schedule::{
17 EnrolmentCheckGroup, EnrolmentCheckSource,
18};
19use headless_lms_utils::services::suotar::SuotarClient;
20use sqlx::PgPool;
21
22use super::commands;
23
24#[derive(Debug, Deserialize)]
25#[serde(rename_all = "camelCase")]
26pub struct RunTickQuery {
27 pub phase: Option<String>,
29 pub course_id: Option<Uuid>,
30 pub course_slug: Option<String>,
32 pub user_id: Option<Uuid>,
33 pub user_email: Option<String>,
34 pub credit_registration_ids: Option<String>,
36 pub account_linking_enabled: Option<bool>,
39}
40
41#[derive(Debug, Serialize, Deserialize, PartialEq)]
42#[serde(rename_all = "camelCase")]
43pub struct UnresolvedScope {
44 pub status: String,
45 pub half: String,
46 pub value: String,
47}
48
49#[derive(Debug, Serialize, Deserialize, PartialEq)]
50#[serde(rename_all = "camelCase", tag = "status")]
51pub enum PhaseTickResult {
52 #[serde(rename_all = "camelCase")]
53 Ran {
54 phase: String,
55 items_processed: i32,
56 items_failed: i32,
57 error: Option<String>,
59 },
60 Skipped { phase: String, reason: String },
62 ScopeNotSupported { phase: String },
64 #[serde(rename_all = "camelCase")]
65 UnknownPhase {
66 phase: Option<String>,
67 known_phases: Vec<String>,
68 },
69}
70
71impl PhaseTickResult {
72 fn of(phase: CreditRegistrationPhase, tick: PhaseTick) -> Self {
73 match tick {
74 PhaseTick::Ran(outcome) => Self::Ran {
75 phase: phase.as_str().to_string(),
76 items_processed: outcome.items_processed,
77 items_failed: outcome.items_failed,
78 error: outcome.error,
79 },
80 PhaseTick::Skipped(reason) => Self::Skipped {
81 phase: phase.as_str().to_string(),
82 reason: match reason {
83 PhaseSkipReason::Paused => "paused".to_string(),
84 PhaseSkipReason::CircuitBreakerOpen => "circuitBreakerOpen".to_string(),
85 PhaseSkipReason::AccountLinkingDisabled => "accountLinkingDisabled".to_string(),
86 PhaseSkipReason::SisuDayGap => "sisuDayGap".to_string(),
87 },
88 },
89 PhaseTick::ScopeNotSupported => Self::ScopeNotSupported {
90 phase: phase.as_str().to_string(),
91 },
92 }
93 }
94}
95
96#[derive(Debug, Serialize, Deserialize, PartialEq)]
98#[serde(rename_all = "camelCase")]
99pub struct RegistrarTickResult {
100 pub phases: Vec<PhaseTickResult>,
101}
102
103async fn run_tick(
107 app_conf: web::Data<ApplicationConfiguration>,
108 pool: web::Data<PgPool>,
109 suotar_client: web::Data<SuotarClient>,
110 query: web::Query<RunTickQuery>,
111) -> ControllerResult<HttpResponse> {
112 super::assert_enabled(&app_conf);
113 let token = skip_authorize();
114
115 let Some(phase) = query
116 .phase
117 .as_deref()
118 .and_then(CreditRegistrationPhase::from_phase_name)
119 else {
120 return token.authorized_ok(HttpResponse::BadRequest().json(
121 PhaseTickResult::UnknownPhase {
122 phase: query.phase.clone(),
123 known_phases: known_phase_names(),
124 },
125 ));
126 };
127
128 let scope = match resolve_scope(&pool, &query).await? {
129 Ok(scope) => scope,
130 Err(unresolved) => {
132 return token.authorized_ok(HttpResponse::BadRequest().json(unresolved));
133 }
134 };
135
136 let ctx = PhaseContext {
137 is_account_linking_enabled: query
138 .account_linking_enabled
139 .unwrap_or(app_conf.suotar_configuration.account_linking_enabled),
140 ..tick_context(&app_conf, &pool, &suotar_client)
141 };
142 debug!(phase = phase.as_str(), ?scope, "run-tick requested");
143 let result = PhaseTickResult::of(phase, run_phase_once(&ctx, phase, &scope).await?);
144 debug!(phase = phase.as_str(), ?result, "run-tick finished");
145 token.authorized_ok(match &result {
146 PhaseTickResult::Ran { .. } | PhaseTickResult::Skipped { .. } => {
147 HttpResponse::Ok().json(&result)
148 }
149 PhaseTickResult::ScopeNotSupported { .. } | PhaseTickResult::UnknownPhase { .. } => {
150 HttpResponse::BadRequest().json(&result)
151 }
152 })
153}
154
155const REGISTRAR_TICK_SEQUENCE: [CreditRegistrationPhase; 5] = [
158 CreditRegistrationPhase::Materialize,
159 CreditRegistrationPhase::Preconditions,
160 CreditRegistrationPhase::ResolveEnrolments,
161 CreditRegistrationPhase::Import,
162 CreditRegistrationPhase::Verify,
163];
164
165async fn run_registrar_tick(
169 app_conf: web::Data<ApplicationConfiguration>,
170 pool: web::Data<PgPool>,
171 suotar_client: web::Data<SuotarClient>,
172) -> ControllerResult<HttpResponse> {
173 super::assert_enabled(&app_conf);
174 let token = skip_authorize();
175
176 let scope = RegistrationScope::default();
177 let ctx = tick_context(&app_conf, &pool, &suotar_client);
178 debug!("run-registrar-tick requested");
179 let mut phases = Vec::new();
180 for phase in REGISTRAR_TICK_SEQUENCE {
181 phases.push(PhaseTickResult::of(
182 phase,
183 run_phase_once(&ctx, phase, &scope).await?,
184 ));
185 }
186 token.authorized_ok(HttpResponse::Ok().json(RegistrarTickResult { phases }))
187}
188
189#[derive(Debug, Deserialize)]
190#[serde(rename_all = "camelCase")]
191pub struct RegradeCompletionPayload {
192 pub credit_registration_id: Uuid,
193 pub grade: Option<i32>,
195 pub passed: Option<bool>,
197}
198
199#[derive(Debug, Serialize, Deserialize, PartialEq)]
200#[serde(rename_all = "camelCase")]
201pub struct RegradeCompletionResult {
202 pub course_module_completion_id: Uuid,
203 pub grade: Option<i32>,
204}
205
206async fn regrade_completion(
212 app_conf: web::Data<ApplicationConfiguration>,
213 pool: web::Data<PgPool>,
214 payload: web::Json<RegradeCompletionPayload>,
215) -> ControllerResult<HttpResponse> {
216 super::assert_enabled(&app_conf);
217 let token = skip_authorize();
218
219 let mut conn = pool.acquire().await?;
220 let registration =
221 models::credit_registrations::get_by_id(&mut conn, payload.credit_registration_id).await?;
222 models::course_module_completions::set_grade_for_testing(
223 &mut conn,
224 registration.course_module_completion_id,
225 payload.grade,
226 payload.passed,
227 )
228 .await?;
229 token.authorized_ok(HttpResponse::Ok().json(RegradeCompletionResult {
230 course_module_completion_id: registration.course_module_completion_id,
231 grade: payload.grade,
232 }))
233}
234
235#[derive(Debug, Deserialize)]
236#[serde(rename_all = "camelCase")]
237pub struct SetTestExclusiveHoldPayload {
238 pub user_email: String,
239 pub course_id: Option<Uuid>,
241 pub hold_secs: i64,
242}
243
244#[derive(Debug, Serialize, Deserialize, PartialEq)]
245#[serde(rename_all = "camelCase")]
246pub struct SetTestExclusiveHoldResult {
247 pub held_until: DateTime<Utc>,
248}
249
250const MAX_TEST_EXCLUSIVE_HOLD_SECS: i64 = 120;
254
255async fn set_test_exclusive_hold(
259 app_conf: web::Data<ApplicationConfiguration>,
260 pool: web::Data<PgPool>,
261 payload: web::Json<SetTestExclusiveHoldPayload>,
262) -> ControllerResult<HttpResponse> {
263 super::assert_enabled(&app_conf);
264 let token = skip_authorize();
265
266 if !(0..=MAX_TEST_EXCLUSIVE_HOLD_SECS).contains(&payload.hold_secs) {
267 return token.authorized_ok(HttpResponse::BadRequest().json(format!(
268 "holdSecs must be between 0 and {MAX_TEST_EXCLUSIVE_HOLD_SECS}"
269 )));
270 }
271
272 let mut conn = pool.acquire().await?;
273 let Some(user_id) = models::user_details::get_active_user_id_by_email_case_insensitive(
274 &mut conn,
275 &payload.user_email,
276 )
277 .await?
278 else {
279 return token.authorized_ok(HttpResponse::BadRequest().json(UnresolvedScope {
280 status: "unresolvedScope".to_string(),
281 half: "userEmail".to_string(),
282 value: payload.user_email.clone(),
283 }));
284 };
285
286 let held_until = Utc::now() + Duration::seconds(payload.hold_secs);
287 models::credit_registrations::testing::set_test_exclusive_hold_for_testing(
288 &mut conn,
289 user_id,
290 payload.course_id,
291 held_until,
292 )
293 .await?;
294
295 token.authorized_ok(HttpResponse::Ok().json(SetTestExclusiveHoldResult { held_until }))
296}
297
298#[derive(Debug, Deserialize)]
299#[serde(rename_all = "camelCase")]
300pub struct ExpireEnrolmentRecheckAllowancePayload {
301 pub credit_registration_id: Uuid,
302 #[serde(default)]
304 pub clear_restarts: bool,
305}
306
307async fn expire_enrolment_recheck_allowance(
310 app_conf: web::Data<ApplicationConfiguration>,
311 pool: web::Data<PgPool>,
312 payload: web::Json<ExpireEnrolmentRecheckAllowancePayload>,
313) -> ControllerResult<HttpResponse> {
314 super::assert_enabled(&app_conf);
315 let token = skip_authorize();
316
317 let mut conn = pool.acquire().await?;
318 models::credit_registrations::testing::expire_enrolment_recheck_allowance_for_testing(
319 &mut conn,
320 payload.credit_registration_id,
321 payload.clear_restarts,
322 )
323 .await?;
324 token.authorized_ok(HttpResponse::Ok().json(()))
325}
326
327#[derive(Debug, Serialize)]
328#[serde(rename_all = "camelCase")]
329pub struct MakeEnrolmentChecksDueResult {
330 pub made_due_count: u64,
331}
332
333async fn make_enrolment_checks_due(
337 app_conf: web::Data<ApplicationConfiguration>,
338 pool: web::Data<PgPool>,
339 query: web::Query<RunTickQuery>,
340) -> ControllerResult<HttpResponse> {
341 super::assert_enabled(&app_conf);
342 let token = skip_authorize();
343
344 let scope = match resolve_scope(&pool, &query).await? {
345 Ok(scope) if !scope.is_unscoped() => scope,
346 Ok(_) => {
347 return token.authorized_ok(
348 HttpResponse::BadRequest().json("A scope is required, or every test's rows move."),
349 );
350 }
351 Err(unresolved) => {
352 return token.authorized_ok(HttpResponse::BadRequest().json(unresolved));
353 }
354 };
355 let mut conn = pool.acquire().await?;
356 let made_due_count =
357 models::credit_registrations::testing::make_enrolment_checks_due_for_testing(
358 &mut conn, &scope,
359 )
360 .await?;
361 token.authorized_ok(HttpResponse::Ok().json(MakeEnrolmentChecksDueResult { made_due_count }))
362}
363
364#[derive(Debug, Serialize)]
365#[serde(rename_all = "camelCase")]
366pub struct MakeRosterListingsDueResult {
367 pub made_due_count: u64,
368}
369
370async fn make_roster_listings_due(
374 app_conf: web::Data<ApplicationConfiguration>,
375 pool: web::Data<PgPool>,
376 query: web::Query<RunTickQuery>,
377) -> ControllerResult<HttpResponse> {
378 super::assert_enabled(&app_conf);
379 let token = skip_authorize();
380
381 let scope = match resolve_scope(&pool, &query).await? {
382 Ok(scope) => scope,
383 Err(unresolved) => {
384 return token.authorized_ok(HttpResponse::BadRequest().json(unresolved));
385 }
386 };
387 let Some(course_id) = scope.course_id else {
388 return token.authorized_ok(
389 HttpResponse::BadRequest().json("A course is required, or every test's codes move."),
390 );
391 };
392 let mut conn = pool.acquire().await?;
393 let made_due_count =
394 models::credit_registration_roster_schedules::testing::make_listings_due_for_testing(
395 &mut conn, course_id,
396 )
397 .await?;
398 registry_health::reset_rate_limits(&scope);
399 token.authorized_ok(HttpResponse::Ok().json(MakeRosterListingsDueResult { made_due_count }))
400}
401
402#[derive(Debug, Deserialize)]
403#[serde(rename_all = "camelCase")]
404pub struct EnrolmentCheckScheduleQuery {
405 pub credit_registration_id: Uuid,
406}
407
408#[derive(Debug, Serialize)]
410#[serde(rename_all = "camelCase")]
411pub struct EnrolmentCheckSchedule {
412 pub state: models::credit_registrations::CreditRegistrationState,
413 pub group: EnrolmentCheckGroup,
414 pub anchor_at: Option<DateTime<Utc>>,
415 pub step: Option<i32>,
416 pub due_at: Option<DateTime<Utc>>,
417 pub next_attempt_at: DateTime<Utc>,
418 pub is_batched: bool,
419 pub source: EnrolmentCheckSource,
420 pub stopped_at: Option<DateTime<Utc>>,
421 pub requested_at: Option<DateTime<Utc>>,
422 pub restart_count: i32,
423 pub checked_at: Option<DateTime<Utc>>,
424 pub first_failed_at: Option<DateTime<Utc>>,
425 pub submit_retry_count: i32,
426 pub error_code: Option<models::credit_registrations::CreditRegistrationErrorCode>,
427 pub seen_enrolment_ids: Option<Vec<String>>,
428 pub claimed_until: Option<DateTime<Utc>>,
429}
430
431async fn enrolment_check_schedule(
432 app_conf: web::Data<ApplicationConfiguration>,
433 pool: web::Data<PgPool>,
434 query: web::Query<EnrolmentCheckScheduleQuery>,
435) -> ControllerResult<HttpResponse> {
436 super::assert_enabled(&app_conf);
437 let token = skip_authorize();
438
439 let mut conn = pool.acquire().await?;
440 let row =
441 models::credit_registrations::get_by_id(&mut conn, query.credit_registration_id).await?;
442 token.authorized_ok(HttpResponse::Ok().json(EnrolmentCheckSchedule {
443 state: row.state,
444 group: row.enrolment_check_group,
445 anchor_at: row.enrolment_check_anchor_at,
446 step: row.enrolment_check_step,
447 due_at: row.enrolment_check_due_at,
448 next_attempt_at: row.next_attempt_at,
449 is_batched: row.is_enrolment_check_batched,
450 source: row.enrolment_check_source,
451 stopped_at: row.enrolment_checks_stopped_at,
452 requested_at: row.enrolment_check_requested_at,
453 restart_count: row.enrolment_check_restart_count,
454 checked_at: row.enrolment_checked_at,
455 first_failed_at: row.first_failed_at,
456 submit_retry_count: row.submit_retry_count,
457 error_code: row.error_code,
458 seen_enrolment_ids: row.seen_enrolment_ids,
459 claimed_until: row.enrolment_check_claimed_until,
460 }))
461}
462
463#[derive(Debug, Deserialize)]
464#[serde(rename_all = "camelCase")]
465pub struct RegisterNewCompletionsViaSuotarPayload {
466 pub course_module_id: Uuid,
467}
468
469async fn register_new_completions_via_suotar(
472 app_conf: web::Data<ApplicationConfiguration>,
473 pool: web::Data<PgPool>,
474 payload: web::Json<RegisterNewCompletionsViaSuotarPayload>,
475) -> ControllerResult<HttpResponse> {
476 super::assert_enabled(&app_conf);
477 let token = skip_authorize();
478
479 let mut conn = pool.acquire().await?;
480 models::course_modules::set_register_eligible_new_completions_via_suotar(
481 &mut conn,
482 payload.course_module_id,
483 true,
484 )
485 .await?;
486 token.authorized_ok(HttpResponse::Ok().json(()))
487}
488
489#[derive(Debug, Deserialize)]
490#[serde(rename_all = "camelCase")]
491pub struct QueuedEmailsQuery {
492 pub user_email: String,
493}
494
495#[derive(Debug, Serialize, Deserialize, PartialEq)]
496#[serde(rename_all = "camelCase")]
497pub struct QueuedEmail {
498 pub template_type: String,
499 pub placeholders: serde_json::Value,
500}
501
502const QUEUED_EMAIL_SCAN: i64 = 200;
506
507async fn queued_emails(
510 app_conf: web::Data<ApplicationConfiguration>,
511 pool: web::Data<PgPool>,
512 query: web::Query<QueuedEmailsQuery>,
513) -> ControllerResult<HttpResponse> {
514 super::assert_enabled(&app_conf);
515 let token = skip_authorize();
516
517 let mut conn = pool.acquire().await?;
518 let Some(user_id) = models::user_details::get_active_user_id_by_email_case_insensitive(
519 &mut conn,
520 &query.user_email,
521 )
522 .await?
523 else {
524 return token.authorized_ok(HttpResponse::BadRequest().json(UnresolvedScope {
525 status: "unresolvedScope".to_string(),
526 half: "userEmail".to_string(),
527 value: query.user_email.clone(),
528 }));
529 };
530 let queued = models::email_deliveries::get_recent_template_types_for_user_for_testing(
531 &mut conn,
532 user_id,
533 QUEUED_EMAIL_SCAN,
534 )
535 .await?
536 .into_iter()
537 .map(|(template_type, placeholders)| QueuedEmail {
538 template_type: serde_json::to_value(template_type)
539 .ok()
540 .and_then(|value| value.as_str().map(str::to_string))
541 .unwrap_or_default(),
542 placeholders,
543 })
544 .collect::<Vec<_>>();
545
546 token.authorized_ok(HttpResponse::Ok().json(queued))
547}
548
549fn tick_context<'a>(
551 app_conf: &'a ApplicationConfiguration,
552 pool: &'a PgPool,
553 suotar_client: &'a SuotarClient,
554) -> PhaseContext<'a> {
555 PhaseContext::from_app(pool, suotar_client, app_conf, Runner::Other("run-tick"))
556}
557
558async fn resolve_scope(
560 pool: &PgPool,
561 query: &RunTickQuery,
562) -> anyhow::Result<Result<RegistrationScope, UnresolvedScope>> {
563 let mut scope = RegistrationScope {
564 course_id: query.course_id,
565 user_id: query.user_id,
566 credit_registration_ids: Vec::new(),
567 };
568 if let Some(slug) = &query.course_slug {
569 let mut conn = pool.acquire().await?;
570 let found = models::courses::get_active_course_id_by_slug(&mut conn, slug).await?;
571 match found {
572 Some(id) => scope.course_id = Some(id),
573 None => {
574 return Ok(Err(UnresolvedScope {
575 status: "unresolvedScope".to_string(),
576 half: "courseSlug".to_string(),
577 value: slug.clone(),
578 }));
579 }
580 }
581 }
582 if let Some(email) = &query.user_email {
583 let mut conn = pool.acquire().await?;
584 let found =
585 models::user_details::get_active_user_id_by_email_case_insensitive(&mut conn, email)
586 .await?;
587 match found {
588 Some(id) => scope.user_id = Some(id),
589 None => {
590 return Ok(Err(UnresolvedScope {
591 status: "unresolvedScope".to_string(),
592 half: "userEmail".to_string(),
593 value: email.clone(),
594 }));
595 }
596 }
597 }
598 if let Some(raw) = &query.credit_registration_ids {
599 for part in raw.split(',').filter(|part| !part.trim().is_empty()) {
600 match Uuid::parse_str(part.trim()) {
601 Ok(id) => scope.credit_registration_ids.push(id),
602 Err(_) => {
603 return Ok(Err(UnresolvedScope {
604 status: "unresolvedScope".to_string(),
605 half: "creditRegistrationIds".to_string(),
606 value: part.trim().to_string(),
607 }));
608 }
609 }
610 }
611 }
612 Ok(Ok(scope))
613}
614
615fn known_phase_names() -> Vec<String> {
616 CreditRegistrationPhase::ALL
617 .iter()
618 .map(|phase| phase.as_str().to_string())
619 .collect()
620}
621
622pub fn _add_routes(cfg: &mut ServiceConfig) {
623 cfg.route("/run-tick", web::post().to(run_tick))
624 .route("/run-registrar-tick", web::post().to(run_registrar_tick))
625 .route("/regrade-completion", web::post().to(regrade_completion))
626 .route(
627 "/test-exclusive-hold",
628 web::post().to(set_test_exclusive_hold),
629 )
630 .route(
631 "/expire-enrolment-recheck-allowance",
632 web::post().to(expire_enrolment_recheck_allowance),
633 )
634 .route(
635 "/enrolment-check-schedule",
636 web::get().to(enrolment_check_schedule),
637 )
638 .route(
639 "/make-enrolment-checks-due",
640 web::post().to(make_enrolment_checks_due),
641 )
642 .route(
643 "/make-roster-listings-due",
644 web::post().to(make_roster_listings_due),
645 )
646 .route(
647 "/register-new-completions-via-suotar",
648 web::post().to(register_new_completions_via_suotar),
649 )
650 .route("/queued-emails", web::get().to(queued_emails))
651 .configure(commands::_add_routes);
652}
653
654#[cfg(test)]
655mod tests {
656 use actix_web::{App, http::StatusCode, test, web::Data};
657
658 use super::*;
659 use crate::controllers::configure_controllers;
660
661 fn app_conf(test_suotar: bool) -> ApplicationConfiguration {
663 ApplicationConfiguration {
664 test_suotar,
665 ..ApplicationConfiguration::mock_conf()
666 .expect("the mock configuration is built from constants")
667 }
668 }
669
670 async fn call_run_tick(
673 test_suotar: bool,
674 query: &str,
675 ) -> actix_web::dev::ServiceResponse<actix_web::body::BoxBody> {
676 let app_conf = Data::new(app_conf(test_suotar));
677 let pool = Data::new(
678 PgPool::connect_lazy("postgres://headless-lms@localhost:54328/headless_lms_dev")
679 .expect("a lazy pool only parses the url"),
680 );
681 let service = test::init_service(
682 App::new()
683 .app_data(pool)
684 .app_data(Data::new(SuotarClient::mock_for_test()))
685 .app_data(app_conf.clone())
686 .service(
687 web::scope("/api/v0")
688 .configure(|cfg| configure_controllers(cfg, app_conf.clone())),
689 ),
690 )
691 .await;
692 let req = test::TestRequest::post()
693 .uri(&format!("/api/v0/mock-suotar/control/run-tick{query}"))
694 .to_request();
695 test::call_service(&service, req).await
696 }
697
698 #[actix_web::test]
700 async fn run_tick_is_absent_when_the_mock_is_disabled() {
701 let res = call_run_tick(false, "?phase=verify").await;
702 assert_eq!(res.status(), StatusCode::NOT_FOUND);
703 }
704
705 #[actix_web::test]
708 async fn a_scope_a_phase_cannot_apply_is_refused() {
709 let res = call_run_tick(
710 true,
711 &format!("?phase=retention-sweep&courseId={}", Uuid::new_v4()),
712 )
713 .await;
714 assert_eq!(res.status(), StatusCode::BAD_REQUEST);
715 let body: PhaseTickResult = test::read_body_json(res).await;
716 assert_eq!(
717 body,
718 PhaseTickResult::ScopeNotSupported {
719 phase: "retention-sweep".to_string()
720 }
721 );
722 }
723
724 #[actix_web::test]
725 async fn an_invented_phase_name_is_rejected() {
726 let res = call_run_tick(true, "?phase=materialise").await;
727 assert_eq!(res.status(), StatusCode::BAD_REQUEST);
728 let body: PhaseTickResult = test::read_body_json(res).await;
729 assert_eq!(
730 body,
731 PhaseTickResult::UnknownPhase {
732 phase: Some("materialise".to_string()),
733 known_phases: known_phase_names(),
734 }
735 );
736 }
737}