1use crate::credit_registration_enrolment_check_outcomes::{self, NewEnrolmentCheckOutcome};
7use crate::credit_registrations::{
8 CreditRegistration, CreditRegistrationState, RegistrationScope, is_waiting_for_enrolment,
9};
10use crate::prelude::*;
11use chrono::TimeDelta;
12
13use super::enrolment_check_schedule::{
14 BATCH_INTERVAL, BATCH_PULL_FORWARD, CHECK_REQUEST_MIN_INTERVAL, CHECK_REQUEST_RESTART_WINDOW,
15 EnrolmentCheckGroup, EnrolmentCheckSource, MAX_CHECK_REQUEST_RESTARTS_PER_DAY,
16 STUDENT_CHECK_REQUEST_MIN_ROW_AGE, ScheduledEnrolmentCheck, TRANSIENT_FAILURE_RETRY,
17 VISIT_RESTART_MIN_INTERVAL, first_check, never, next_check_after,
18};
19use super::study_registry::RegistryEnrolment;
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
23pub enum CheckRequestOutcome {
24 CheckStarted,
26 Rescheduled,
28 TooSoon,
31 NotWaiting,
33}
34
35impl CheckRequestOutcome {
36 pub fn started_check(self) -> bool {
38 self == Self::CheckStarted
39 }
40}
41
42struct ScheduleFacts {
43 state: CreditRegistrationState,
44 no_usable_enrolment_since: Option<DateTime<Utc>>,
45 enrolment_check_group: EnrolmentCheckGroup,
46 enrolment_check_anchor_at: Option<DateTime<Utc>>,
47 enrolment_checks_stopped_at: Option<DateTime<Utc>>,
48 enrolment_check_requested_at: Option<DateTime<Utc>>,
49 enrolment_check_restart_window_started_at: Option<DateTime<Utc>>,
50 enrolment_check_restart_count: i32,
51 enrolment_checked_at: Option<DateTime<Utc>>,
52}
53
54async fn lock_schedule(conn: &mut PgConnection, id: Uuid) -> ModelResult<ScheduleFacts> {
55 let res = sqlx::query_as!(
56 ScheduleFacts,
57 r#"
58SELECT state,
59 no_usable_enrolment_since,
60 enrolment_check_group,
61 enrolment_check_anchor_at,
62 enrolment_checks_stopped_at,
63 enrolment_check_requested_at,
64 enrolment_check_restart_window_started_at,
65 enrolment_check_restart_count,
66 enrolment_checked_at
67FROM credit_registrations
68WHERE id = $1
69 AND deleted_at IS NULL
70FOR UPDATE
71 "#,
72 id,
73 )
74 .fetch_one(conn)
75 .await?;
76 Ok(res)
77}
78
79impl ScheduleFacts {
80 fn is_waiting(&self) -> bool {
81 is_waiting_for_enrolment(
82 self.state,
83 self.enrolment_check_anchor_at,
84 self.no_usable_enrolment_since,
85 )
86 }
87}
88
89fn is_within(at: Option<DateTime<Utc>>, now: DateTime<Utc>, span: TimeDelta) -> bool {
90 at.is_some_and(|at| now - at < span)
91}
92
93#[allow(clippy::too_many_arguments)]
97async fn write_schedule(
98 conn: &mut PgConnection,
99 id: Uuid,
100 group: EnrolmentCheckGroup,
101 anchor_at: DateTime<Utc>,
102 scheduled: Option<ScheduledEnrolmentCheck>,
103 source: EnrolmentCheckSource,
104 next_attempt_at: DateTime<Utc>,
105 replaces_next_attempt: bool,
106) -> ModelResult<()> {
107 sqlx::query!(
108 r#"
109UPDATE credit_registrations
110SET enrolment_check_group = GREATEST(enrolment_check_group, $2),
111 enrolment_check_anchor_at = $3,
112 enrolment_check_step = $4,
113 enrolment_check_due_at = $5,
114 is_enrolment_check_batched = $6,
115 enrolment_checks_stopped_at = CASE
116 WHEN $4::int IS NULL THEN COALESCE(enrolment_checks_stopped_at, now())
117 END,
118 enrolment_check_source = $7,
119 next_attempt_at = CASE
120 WHEN state <> 'no_usable_enrolment' THEN next_attempt_at
121 WHEN $9 THEN $8
122 ELSE LEAST(next_attempt_at, $8)
123 END
124WHERE id = $1
125 "#,
126 id,
127 group as EnrolmentCheckGroup,
128 anchor_at,
129 scheduled.map(|scheduled| scheduled.step),
130 scheduled.map(|scheduled| scheduled.due_at),
131 scheduled.is_some_and(|scheduled| scheduled.is_batched),
132 source as EnrolmentCheckSource,
133 next_attempt_at,
134 replaces_next_attempt,
135 )
136 .execute(conn)
137 .await?;
138 Ok(())
139}
140
141pub async fn request_check(
148 conn: &mut PgConnection,
149 id: Uuid,
150 source: EnrolmentCheckSource,
151 now: DateTime<Utc>,
152) -> ModelResult<CheckRequestOutcome> {
153 let mut tx = conn.begin().await?;
154 let facts = lock_schedule(&mut tx, id).await?;
155 if !facts.is_waiting() {
156 return Ok(CheckRequestOutcome::NotWaiting);
157 }
158 if is_within(
159 facts.enrolment_check_requested_at,
160 now,
161 CHECK_REQUEST_MIN_INTERVAL,
162 ) {
163 return Ok(CheckRequestOutcome::TooSoon);
164 }
165 let checked_recently = is_within(facts.enrolment_checked_at, now, CHECK_REQUEST_MIN_INTERVAL);
166 let window_is_current = is_within(
167 facts.enrolment_check_restart_window_started_at,
168 now,
169 CHECK_REQUEST_RESTART_WINDOW,
170 );
171 let (window_started_at, restart_count) = if window_is_current {
172 (
173 facts
174 .enrolment_check_restart_window_started_at
175 .unwrap_or(now),
176 facts.enrolment_check_restart_count,
177 )
178 } else {
179 (now, 0)
180 };
181 let may_restart = restart_count < MAX_CHECK_REQUEST_RESTARTS_PER_DAY;
182 if !may_restart && checked_recently {
183 return Ok(CheckRequestOutcome::TooSoon);
184 }
185 sqlx::query!(
186 r#"
187UPDATE credit_registrations
188SET enrolment_check_requested_at = $2,
189 enrolment_check_restart_window_started_at = $3,
190 enrolment_check_restart_count = $4
191WHERE id = $1
192 "#,
193 id,
194 now,
195 window_started_at,
196 restart_count + i32::from(may_restart),
197 )
198 .execute(&mut *tx)
199 .await?;
200
201 let outcome = if may_restart {
202 let group = facts
203 .enrolment_check_group
204 .max(EnrolmentCheckGroup::CheckRequested);
205 let (scheduled, source, next_attempt_at, outcome) = if checked_recently {
206 let scheduled = next_check_after(group, now, now);
207 (
208 scheduled,
209 EnrolmentCheckSource::Schedule,
210 scheduled.map_or_else(never, |scheduled| scheduled.release_at()),
211 CheckRequestOutcome::Rescheduled,
212 )
213 } else {
214 (
215 first_check(group, now),
216 source,
217 now,
218 CheckRequestOutcome::CheckStarted,
219 )
220 };
221 write_schedule(
222 &mut tx,
223 id,
224 group,
225 now,
226 scheduled,
227 source,
228 next_attempt_at,
229 true,
230 )
231 .await?;
232 outcome
233 } else {
234 sqlx::query!(
235 r#"
236UPDATE credit_registrations
237SET enrolment_check_source = $2,
238 next_attempt_at = CASE
239 WHEN state = 'no_usable_enrolment' THEN now()
240 ELSE next_attempt_at
241 END
242WHERE id = $1
243 "#,
244 id,
245 source as EnrolmentCheckSource,
246 )
247 .execute(&mut *tx)
248 .await?;
249 CheckRequestOutcome::CheckStarted
250 };
251 tx.commit().await?;
252 Ok(outcome)
253}
254
255pub fn is_check_request_limited(
258 enrolment_check_requested_at: Option<DateTime<Utc>>,
259 enrolment_checked_at: Option<DateTime<Utc>>,
260 now: DateTime<Utc>,
261) -> bool {
262 is_within(
263 enrolment_check_requested_at,
264 now,
265 CHECK_REQUEST_MIN_INTERVAL,
266 ) || is_within(enrolment_checked_at, now, CHECK_REQUEST_MIN_INTERVAL)
267}
268
269pub fn is_too_new_for_student_check_request(
272 row_created_at: DateTime<Utc>,
273 now: DateTime<Utc>,
274) -> bool {
275 now - row_created_at < STUDENT_CHECK_REQUEST_MIN_ROW_AGE
276}
277
278pub async fn record_visit(
282 conn: &mut PgConnection,
283 id: Uuid,
284 now: DateTime<Utc>,
285) -> ModelResult<bool> {
286 let mut tx = conn.begin().await?;
287 let facts = lock_schedule(&mut tx, id).await?;
288 if !facts.is_waiting() {
289 return Ok(false);
290 }
291 let restart_is_due = !is_within(
292 facts.enrolment_check_anchor_at,
293 now,
294 VISIT_RESTART_MIN_INTERVAL,
295 );
296 let restarts = facts.enrolment_check_group < EnrolmentCheckGroup::Visited
297 || (restart_is_due
298 && (facts.enrolment_checks_stopped_at.is_some()
299 || facts.enrolment_check_group == EnrolmentCheckGroup::Visited));
300 if !restarts {
301 return Ok(false);
302 }
303 let group = facts
304 .enrolment_check_group
305 .max(EnrolmentCheckGroup::Visited);
306 let scheduled = first_check(group, now);
307 write_schedule(
308 &mut tx,
309 id,
310 group,
311 now,
312 scheduled,
313 EnrolmentCheckSource::Schedule,
314 scheduled.map_or_else(never, |scheduled| scheduled.release_at()),
315 false,
316 )
317 .await?;
318 tx.commit().await?;
319 Ok(true)
320}
321
322pub async fn schedule_next_check(conn: &mut PgConnection, id: Uuid) -> ModelResult<()> {
329 let facts = lock_schedule(conn, id).await?;
330 if facts.state != CreditRegistrationState::NoUsableEnrolment {
331 return Ok(());
332 }
333 let now = Utc::now();
334 let anchor_at = facts.enrolment_check_anchor_at.unwrap_or(now);
335 let scheduled = next_check_after(facts.enrolment_check_group, anchor_at, now);
336 write_schedule(
337 conn,
338 id,
339 facts.enrolment_check_group,
340 anchor_at,
341 scheduled,
342 EnrolmentCheckSource::Schedule,
343 scheduled.map_or_else(never, |scheduled| scheduled.release_at()),
344 true,
345 )
346 .await
347}
348
349#[derive(Debug, Clone, PartialEq)]
352pub struct EnrolmentCheckStart {
353 pub credit_registration_id: Uuid,
354 pub group: EnrolmentCheckGroup,
355 pub anchor_at: DateTime<Utc>,
356 pub scheduled: ScheduledEnrolmentCheck,
357 pub source: EnrolmentCheckSource,
358}
359
360pub async fn record_starts(
366 conn: &mut PgConnection,
367 starts: &[EnrolmentCheckStart],
368) -> ModelResult<()> {
369 if starts.is_empty() {
370 return Ok(());
371 }
372 let ids: Vec<Uuid> = starts
373 .iter()
374 .map(|start| start.credit_registration_id)
375 .collect();
376 let groups: Vec<EnrolmentCheckGroup> = starts.iter().map(|start| start.group).collect();
377 let anchors: Vec<DateTime<Utc>> = starts.iter().map(|start| start.anchor_at).collect();
378 let steps: Vec<i32> = starts.iter().map(|start| start.scheduled.step).collect();
379 let due_ats: Vec<DateTime<Utc>> = starts.iter().map(|start| start.scheduled.due_at).collect();
380 let batched: Vec<bool> = starts
381 .iter()
382 .map(|start| start.scheduled.is_batched)
383 .collect();
384 let sources: Vec<EnrolmentCheckSource> = starts.iter().map(|start| start.source).collect();
385 sqlx::query!(
386 r#"
387UPDATE credit_registrations cr
388SET enrolment_check_group = GREATEST(cr.enrolment_check_group, started.enrolment_check_group),
389 enrolment_check_anchor_at = started.anchor_at,
390 enrolment_check_step = started.step,
391 enrolment_check_due_at = started.due_at,
392 is_enrolment_check_batched = started.is_batched,
393 enrolment_checks_stopped_at = NULL,
394 enrolment_check_source = started.source
395FROM UNNEST(
396 $1::uuid [],
397 $2::enrolment_check_group [],
398 $3::timestamptz [],
399 $4::int [],
400 $5::timestamptz [],
401 $6::boolean [],
402 $7::enrolment_check_source []
403 ) AS started(
404 id,
405 enrolment_check_group,
406 anchor_at,
407 step,
408 due_at,
409 is_batched,
410 source
411 )
412WHERE cr.id = started.id
413 AND cr.state IN ('no_usable_enrolment', 'ready_to_submit')
414 AND cr.deleted_at IS NULL
415 "#,
416 &ids,
417 &groups as &[EnrolmentCheckGroup],
418 &anchors,
419 &steps,
420 &due_ats,
421 &batched,
422 &sources as &[EnrolmentCheckSource],
423 )
424 .execute(&mut *conn)
425 .await?;
426 sqlx::query!(
427 r#"
428UPDATE credit_registrations replaced
429SET enrolment_check_anchor_at = COALESCE(replaced.enrolment_check_anchor_at, now()),
430 enrolment_check_step = NULL,
431 enrolment_check_due_at = NULL,
432 is_enrolment_check_batched = FALSE,
433 enrolment_checks_stopped_at = COALESCE(replaced.enrolment_checks_stopped_at, now()),
434 next_attempt_at = $2
435FROM credit_registrations started
436WHERE started.id = ANY($1::uuid [])
437 AND replaced.user_id = started.user_id
438 AND replaced.course_module_id = started.course_module_id
439 AND replaced.id <> started.id
440 AND replaced.created_at < started.created_at
441 AND replaced.state = 'no_usable_enrolment'
442 AND replaced.superseded_by_id IS NULL
443 AND replaced.deleted_at IS NULL
444 "#,
445 &ids,
446 never(),
447 )
448 .execute(conn)
449 .await?;
450 Ok(())
451}
452
453#[derive(Debug, Clone, PartialEq)]
455pub struct RosterEnrolee {
456 pub user_id: Uuid,
457 pub enrolment_ids: Vec<String>,
458}
459
460pub async fn wake_for_roster_listing(
467 conn: &mut PgConnection,
468 course_module_id: Uuid,
469 enrolees: &[RosterEnrolee],
470) -> ModelResult<u64> {
471 let mut user_ids = Vec::new();
472 let mut enrolment_ids: Vec<Option<String>> = Vec::new();
473 for enrolee in enrolees {
474 if enrolee.enrolment_ids.is_empty() {
475 user_ids.push(enrolee.user_id);
476 enrolment_ids.push(None);
477 }
478 for enrolment_id in &enrolee.enrolment_ids {
479 user_ids.push(enrolee.user_id);
480 enrolment_ids.push(Some(enrolment_id.clone()));
481 }
482 }
483 let res = sqlx::query!(
484 r#"
485WITH listed AS (
486 SELECT user_id,
487 ARRAY_REMOVE(ARRAY_AGG(enrolment_id), NULL) AS enrolment_ids
488 FROM UNNEST($2::uuid [], $3::text []) AS listing(user_id, enrolment_id)
489 GROUP BY user_id
490)
491UPDATE credit_registrations cr
492SET enrolment_check_source = CASE
493 WHEN cr.next_attempt_at > now() THEN 'roster_listing'
494 ELSE cr.enrolment_check_source
495 END,
496 next_attempt_at = LEAST(cr.next_attempt_at, now()),
497 seen_enrolment_ids = ARRAY(
498 SELECT DISTINCT UNNEST(COALESCE(cr.seen_enrolment_ids, '{}') || listed.enrolment_ids)
499 )
500FROM listed
501WHERE cr.course_module_id = $1
502 AND cr.user_id = listed.user_id
503 AND cr.state = 'no_usable_enrolment'
504 AND cr.superseded_by_id IS NULL
505 AND cr.deleted_at IS NULL
506 AND (
507 cr.seen_enrolment_ids IS NULL
508 OR NOT listed.enrolment_ids <@ cr.seen_enrolment_ids
509 )
510 AND NOT EXISTS (
511 SELECT 1
512 FROM credit_registrations later
513 WHERE later.user_id = cr.user_id
514 AND later.course_module_id = cr.course_module_id
515 AND later.created_at > cr.created_at
516 AND later.deleted_at IS NULL
517 )
518 "#,
519 course_module_id,
520 &user_ids,
521 &enrolment_ids as &[Option<String>],
522 )
523 .execute(conn)
524 .await?;
525 Ok(res.rows_affected())
526}
527
528async fn add_seen_enrolment_ids(
531 conn: &mut PgConnection,
532 id: Uuid,
533 enrolment_ids: &[String],
534) -> ModelResult<()> {
535 sqlx::query!(
536 r#"
537UPDATE credit_registrations
538SET seen_enrolment_ids = ARRAY(
539 SELECT DISTINCT UNNEST(COALESCE(seen_enrolment_ids, '{}') || $2::text [])
540 )
541WHERE id = $1
542 "#,
543 id,
544 enrolment_ids,
545 )
546 .execute(conn)
547 .await?;
548 Ok(())
549}
550
551#[derive(Clone, Copy)]
553pub struct EnrolmentCheckAnswer<'a> {
554 pub checked: &'a CreditRegistration,
556 pub usable_enrolment: Option<&'a RegistryEnrolment>,
558 pub listed_enrolments: &'a [RegistryEnrolment],
559}
560
561pub async fn record_enrolment_check(
565 conn: &mut PgConnection,
566 check: Option<&EnrolmentCheckAnswer<'_>>,
567 after: &CreditRegistration,
568) -> ModelResult<()> {
569 let Some(check) = check else {
570 return Ok(());
571 };
572 let row = check.checked;
573 let was_answered = after.state != CreditRegistrationState::NoUsableEnrolment
574 || after.enrolment_checked_at != row.enrolment_checked_at;
575 if !was_answered {
576 return Ok(());
577 }
578 credit_registration_enrolment_check_outcomes::insert(
579 conn,
580 &NewEnrolmentCheckOutcome {
581 credit_registration_id: row.id,
582 course_module_id: row.course_module_id,
583 enrolment_check_group: row.enrolment_check_group,
584 enrolment_check_step: row.enrolment_check_step,
585 source: row.enrolment_check_source,
586 due_at: row.enrolment_check_due_at,
587 checked_at: after.enrolment_checked_at.unwrap_or_else(Utc::now),
588 previous_checked_at: row.enrolment_checked_at,
589 is_enrolment_found: check.usable_enrolment.is_some(),
590 enrolled_at: check
591 .usable_enrolment
592 .and_then(|enrolment| enrolment.enrolment_date_time),
593 },
594 )
595 .await?;
596 let seen: Vec<String> = check
597 .listed_enrolments
598 .iter()
599 .map(|enrolment| enrolment.id.clone())
600 .collect();
601 add_seen_enrolment_ids(conn, row.id, &seen).await?;
602 Ok(())
603}
604
605pub async fn shift_past_pause(
608 conn: &mut PgConnection,
609 course_module_id: Uuid,
610 paused_for: TimeDelta,
611) -> ModelResult<u64> {
612 let res = sqlx::query!(
613 r#"
614UPDATE credit_registrations
615SET enrolment_check_anchor_at = enrolment_check_anchor_at + $2::interval,
616 enrolment_check_due_at = enrolment_check_due_at + $2::interval,
617 next_attempt_at = CASE
618 WHEN state <> 'no_usable_enrolment'
619 OR enrolment_check_source <> 'schedule' THEN next_attempt_at
620 WHEN is_enrolment_check_batched THEN to_timestamp(
621 CEIL(
622 EXTRACT(
623 EPOCH
624 FROM next_attempt_at + $2::interval
625 ) / $3::bigint
626 ) * $3::bigint
627 )
628 ELSE next_attempt_at + $2::interval
629 END
630WHERE course_module_id = $1
631 AND enrolment_check_anchor_at IS NOT NULL
632 AND enrolment_checks_stopped_at IS NULL
633 AND superseded_by_id IS NULL
634 AND deleted_at IS NULL
635 "#,
636 course_module_id,
637 paused_for.max(TimeDelta::zero()) as TimeDelta,
638 BATCH_INTERVAL.num_seconds(),
639 )
640 .execute(conn)
641 .await?;
642 Ok(res.rows_affected())
643}
644
645pub async fn pull_forward_batched_checks(
655 conn: &mut PgConnection,
656 scope: &RegistrationScope,
657) -> ModelResult<()> {
658 sqlx::query!(
659 r#"
660WITH scoped AS (
661 SELECT cr.id,
662 cr.next_attempt_at,
663 cr.enrolment_check_due_at,
664 cr.last_attempt_at,
665 cr.enrolment_check_claimed_until
666 FROM credit_registrations cr
667 JOIN credit_registration_active_course_modules acm ON acm.course_module_id = cr.course_module_id
668 WHERE cr.state = 'no_usable_enrolment'
669 AND cr.is_enrolment_check_batched
670 AND cr.superseded_by_id IS NULL
671 AND cr.deleted_at IS NULL
672 AND ($1::uuid IS NULL OR cr.course_id = $1)
673 AND ($2::uuid IS NULL OR cr.user_id = $2)
674 AND (
675 cardinality($3::uuid []) = 0
676 OR cr.id = ANY($3::uuid [])
677 )
678)
679UPDATE credit_registrations cr
680SET next_attempt_at = now()
681FROM scoped
682WHERE cr.id = scoped.id
683 AND scoped.next_attempt_at > now()
684 AND scoped.enrolment_check_due_at > now()
685 AND scoped.enrolment_check_due_at <= now() + $4::interval
686 AND (
687 scoped.last_attempt_at IS NULL
688 OR scoped.last_attempt_at <= now() - $5::interval
689 )
690 AND EXISTS (
691 SELECT 1
692 FROM scoped released
693 WHERE released.next_attempt_at <= now()
694 AND (
695 released.enrolment_check_claimed_until IS NULL
696 OR released.enrolment_check_claimed_until <= now()
697 )
698 )
699 "#,
700 scope.course_id,
701 scope.user_id,
702 &scope.credit_registration_ids,
703 BATCH_PULL_FORWARD as TimeDelta,
704 TRANSIENT_FAILURE_RETRY as TimeDelta,
705 )
706 .execute(conn)
707 .await?;
708 Ok(())
709}