1use std::collections::HashMap;
8
9use chrono::NaiveDate;
10use utoipa::ToSchema;
11
12use crate::course_module_suotar_configurations::ModuleToList;
13use crate::credit_registrations::CreditRegistrationErrorCode;
14use crate::prelude::*;
15
16const HOUR_SECS: i64 = 60 * 60;
17const DAY_SECS: i64 = 24 * HOUR_SECS;
18
19pub const ACTIVE_WINDOW_SECS: i64 = 7 * DAY_SECS;
21pub const DORMANT_AFTER_SECS: i64 = 60 * DAY_SECS;
24pub const ACTIVE_INTERVAL_SECS: i64 = 8 * HOUR_SECS;
25pub const IDLE_INTERVAL_SECS: i64 = DAY_SECS;
26pub const DORMANT_INTERVAL_SECS: i64 = 7 * DAY_SECS;
27pub const TRIGGER_MIN_GAP_SECS: i64 = HOUR_SECS;
30pub const VISIT_FOLLOW_UP_SECS: i64 = 2 * HOUR_SECS;
32pub const MAX_TRIGGERED_FETCHES_PER_DAY: i32 = 6;
33const ALONE_FAILURE_BACKOFF_SECS: [i64; 3] = [HOUR_SECS, 4 * HOUR_SECS, DAY_SECS];
35
36#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
38#[serde(rename_all = "snake_case")]
39pub enum RosterTier {
40 Active,
42 Idle,
44 Dormant,
46 Unlisted,
49}
50
51impl RosterTier {
52 pub fn interval_secs(self) -> Option<i64> {
55 match self {
56 Self::Active => Some(ACTIVE_INTERVAL_SECS),
57 Self::Idle => Some(IDLE_INTERVAL_SECS),
58 Self::Dormant => Some(DORMANT_INTERVAL_SECS),
59 Self::Unlisted => None,
60 }
61 }
62}
63
64#[derive(Debug, Clone, PartialEq)]
66pub struct RosterTierFacts {
67 pub last_completion_at: Option<DateTime<Utc>>,
70 pub has_waiting_rows: bool,
71 pub window_closed_at: Option<DateTime<Utc>>,
72}
73
74pub fn roster_tier(
76 facts: &RosterTierFacts,
77 is_account_linking_enabled: bool,
78 now: DateTime<Utc>,
79) -> RosterTier {
80 let is_window_closed = facts.window_closed_at.is_some_and(|closed| {
81 facts
82 .last_completion_at
83 .is_none_or(|completion| completion <= closed)
84 });
85 if is_window_closed || (!is_account_linking_enabled && !facts.has_waiting_rows) {
86 return RosterTier::Unlisted;
87 }
88 let quiet_secs = facts
89 .last_completion_at
90 .map_or(i64::MAX, |completion| (now - completion).num_seconds());
91 if quiet_secs < ACTIVE_WINDOW_SECS {
92 RosterTier::Active
93 } else if quiet_secs < DORMANT_AFTER_SECS {
94 RosterTier::Idle
95 } else if is_account_linking_enabled {
96 RosterTier::Dormant
97 } else {
98 RosterTier::Unlisted
99 }
100}
101
102#[derive(Debug, Clone, PartialEq)]
104pub struct RosterSchedule {
105 pub course_code: String,
106 pub last_fetched_at: Option<DateTime<Utc>>,
107 pub last_fetch_duration_ms: Option<i32>,
108 pub last_listed_person_count: Option<i32>,
109 pub triggered_fetch_at: Option<DateTime<Utc>>,
110 pub follow_up_fetch_at: Option<DateTime<Utc>>,
111 pub triggered_fetch_day: Option<NaiveDate>,
112 pub triggered_fetch_count: i32,
113 pub is_fetched_alone: bool,
114 pub consecutive_failures: i32,
115 pub retry_not_before: Option<DateTime<Utc>>,
116 pub last_error: Option<CreditRegistrationErrorCode>,
117 pub module_count: i64,
119 pub tier_facts: RosterTierFacts,
120}
121
122impl RosterSchedule {
123 pub fn next_fetch_at(
125 &self,
126 is_account_linking_enabled: bool,
127 now: DateTime<Utc>,
128 ) -> Option<DateTime<Utc>> {
129 let tier_due = roster_tier(&self.tier_facts, is_account_linking_enabled, now)
130 .interval_secs()
131 .map(|interval| {
132 self.last_fetched_at
133 .map_or(now, |fetched| fetched + chrono::Duration::seconds(interval))
134 });
135 [tier_due, self.triggered_due_at()]
136 .into_iter()
137 .flatten()
138 .min()
139 }
140
141 pub fn triggered_due_at(&self) -> Option<DateTime<Utc>> {
143 [self.triggered_fetch_at, self.follow_up_fetch_at]
144 .into_iter()
145 .flatten()
146 .min()
147 }
148
149 pub fn is_due(&self, is_account_linking_enabled: bool, now: DateTime<Utc>) -> bool {
151 self.retry_not_before.is_none_or(|retry| retry <= now)
152 && self
153 .next_fetch_at(is_account_linking_enabled, now)
154 .is_some_and(|due| due <= now)
155 }
156
157 pub fn is_triggered_due(&self, now: DateTime<Utc>) -> bool {
159 self.triggered_due_at().is_some_and(|due| due <= now)
160 }
161}
162
163pub async fn ensure_rows(conn: &mut PgConnection, course_id: Option<Uuid>) -> ModelResult<()> {
166 sqlx::query!(
167 r#"
168INSERT INTO course_module_suotar_configurations (course_module_id)
169SELECT acm.course_module_id
170FROM credit_registration_active_course_modules acm
171 JOIN course_modules cm ON cm.id = acm.course_module_id AND cm.deleted_at IS NULL
172WHERE TRIM(COALESCE(cm.uh_course_code, '')) <> ''
173 AND ($1::uuid IS NULL OR acm.course_id = $1) ON CONFLICT (course_module_id) DO NOTHING
174 "#,
175 course_id,
176 )
177 .execute(&mut *conn)
178 .await?;
179 sqlx::query!(
180 r#"
181INSERT INTO credit_registration_roster_schedules (course_code)
182SELECT DISTINCT TRIM(cm.uh_course_code)
183FROM credit_registration_active_course_modules acm
184 JOIN course_modules cm ON cm.id = acm.course_module_id AND cm.deleted_at IS NULL
185WHERE TRIM(COALESCE(cm.uh_course_code, '')) <> ''
186 AND ($1::uuid IS NULL OR acm.course_id = $1) ON CONFLICT (course_code) DO NOTHING
187 "#,
188 course_id,
189 )
190 .execute(conn)
191 .await?;
192 Ok(())
193}
194
195#[derive(Debug, Clone, Copy, PartialEq, Eq)]
197pub enum ScheduleSelection {
198 Every,
199 DueCandidates,
202}
203
204pub async fn get_schedules(
207 conn: &mut PgConnection,
208 course_id: Option<Uuid>,
209 selection: ScheduleSelection,
210) -> ModelResult<Vec<RosterSchedule>> {
211 let rows = sqlx::query!(
212 r#"
213WITH modules AS (
214 SELECT acm.course_module_id,
215 TRIM(cm.uh_course_code) AS course_code
216 FROM credit_registration_active_course_modules acm
217 JOIN course_modules cm ON cm.id = acm.course_module_id AND cm.deleted_at IS NULL
218 WHERE TRIM(COALESCE(cm.uh_course_code, '')) <> ''
219 AND ($1::uuid IS NULL OR acm.course_id = $1)
220),
221candidates AS (
222 SELECT s.*
223 FROM credit_registration_roster_schedules s
224 WHERE s.course_code IN (
225 SELECT course_code
226 FROM modules
227 )
228 AND (
229 NOT $2::boolean
230 OR (
231 (
232 s.retry_not_before IS NULL
233 OR s.retry_not_before <= now()
234 )
235 AND (
236 s.last_fetched_at IS NULL
237 OR s.last_fetched_at <= now() - ($3::bigint * INTERVAL '1 second')
238 OR LEAST(s.triggered_fetch_at, s.follow_up_fetch_at) <= now()
239 )
240 )
241 )
242)
243SELECT c.course_code,
244 c.last_fetched_at,
245 c.last_fetch_duration_ms,
246 c.last_listed_person_count,
247 c.triggered_fetch_at,
248 c.follow_up_fetch_at,
249 c.triggered_fetch_day,
250 c.triggered_fetch_count,
251 c.is_fetched_alone,
252 c.consecutive_failures,
253 c.retry_not_before,
254 c.last_error,
255 c.window_closed_at,
256 COUNT(*) AS "module_count!",
257 MAX(facts.last_completion_at) AS last_completion_at,
258 bool_or(facts.has_waiting_rows) AS "has_waiting_rows!"
259FROM candidates c
260 JOIN modules m ON m.course_code = c.course_code
261 CROSS JOIN LATERAL (
262 SELECT (
263 SELECT cr.created_at
264 FROM credit_registrations cr
265 WHERE cr.course_module_id = m.course_module_id
266 AND cr.deleted_at IS NULL
267 ORDER BY cr.created_at DESC
268 LIMIT 1
269 ) AS last_completion_at,
270 EXISTS (
271 SELECT 1
272 FROM credit_registrations cr
273 WHERE cr.course_module_id = m.course_module_id
274 AND cr.state = 'no_usable_enrolment'
275 AND cr.enrolment_checks_stopped_at IS NULL
276 AND cr.superseded_by_id IS NULL
277 AND cr.deleted_at IS NULL
278 ) AS has_waiting_rows
279 ) facts
280GROUP BY c.id,
281 c.course_code,
282 c.last_fetched_at,
283 c.last_fetch_duration_ms,
284 c.last_listed_person_count,
285 c.triggered_fetch_at,
286 c.follow_up_fetch_at,
287 c.triggered_fetch_day,
288 c.triggered_fetch_count,
289 c.is_fetched_alone,
290 c.consecutive_failures,
291 c.retry_not_before,
292 c.last_error,
293 c.window_closed_at
294ORDER BY c.course_code
295 "#,
296 course_id,
297 selection == ScheduleSelection::DueCandidates,
298 ACTIVE_INTERVAL_SECS,
299 )
300 .fetch_all(conn)
301 .await?;
302 Ok(rows
303 .into_iter()
304 .map(|row| RosterSchedule {
305 course_code: row.course_code,
306 last_fetched_at: row.last_fetched_at,
307 last_fetch_duration_ms: row.last_fetch_duration_ms,
308 last_listed_person_count: row.last_listed_person_count,
309 triggered_fetch_at: row.triggered_fetch_at,
310 follow_up_fetch_at: row.follow_up_fetch_at,
311 triggered_fetch_day: row.triggered_fetch_day,
312 triggered_fetch_count: row.triggered_fetch_count,
313 is_fetched_alone: row.is_fetched_alone,
314 consecutive_failures: row.consecutive_failures,
315 retry_not_before: row.retry_not_before,
316 last_error: row.last_error,
317 module_count: row.module_count,
318 tier_facts: RosterTierFacts {
319 last_completion_at: row.last_completion_at,
320 has_waiting_rows: row.has_waiting_rows,
321 window_closed_at: row.window_closed_at,
322 },
323 })
324 .collect())
325}
326
327pub async fn get_modules_by_code(
330 conn: &mut PgConnection,
331 course_id: Option<Uuid>,
332 course_codes: &[String],
333) -> ModelResult<HashMap<String, Vec<ModuleToList>>> {
334 let modules = sqlx::query_as!(
335 ModuleToList,
336 r#"
337SELECT acm.course_module_id AS "course_module_id!",
338 acm.course_id AS "course_id!",
339 TRIM(cm.uh_course_code) AS "uh_course_code!",
340 co.language_code AS "course_language_code!"
341FROM credit_registration_active_course_modules acm
342 JOIN course_modules cm ON cm.id = acm.course_module_id AND cm.deleted_at IS NULL
343 JOIN courses co ON co.id = acm.course_id AND co.deleted_at IS NULL
344WHERE TRIM(cm.uh_course_code) = ANY($2::text [])
345 AND ($1::uuid IS NULL OR acm.course_id = $1)
346ORDER BY cm.id
347 "#,
348 course_id,
349 course_codes,
350 )
351 .fetch_all(conn)
352 .await?;
353 let mut modules_by_code: HashMap<String, Vec<ModuleToList>> = HashMap::new();
354 for module in modules {
355 modules_by_code
356 .entry(module.uh_course_code.clone())
357 .or_default()
358 .push(module);
359 }
360 Ok(modules_by_code)
361}
362
363pub async fn mark_attempted(conn: &mut PgConnection, course_codes: &[String]) -> ModelResult<()> {
365 sqlx::query!(
366 r#"
367UPDATE credit_registration_roster_schedules
368SET last_attempted_at = now()
369WHERE course_code = ANY($1::text [])
370 "#,
371 course_codes,
372 )
373 .execute(conn)
374 .await?;
375 Ok(())
376}
377
378pub async fn mark_fetched(
381 conn: &mut PgConnection,
382 course_code: &str,
383 listed_person_count: i32,
384 duration_ms: i32,
385) -> ModelResult<()> {
386 sqlx::query!(
387 r#"
388UPDATE credit_registration_roster_schedules
389SET last_fetched_at = now(),
390 last_fetch_duration_ms = $3,
391 last_listed_person_count = $2,
392 is_fetched_alone = FALSE,
393 consecutive_failures = 0,
394 retry_not_before = NULL,
395 last_error = NULL,
396 window_closed_at = NULL,
397 triggered_fetch_count = CASE
398 WHEN LEAST(triggered_fetch_at, follow_up_fetch_at) > now()
399 OR COALESCE(triggered_fetch_at, follow_up_fetch_at) IS NULL THEN triggered_fetch_count
400 WHEN triggered_fetch_day = (now() AT TIME ZONE 'UTC')::date THEN triggered_fetch_count + 1
401 ELSE 1
402 END,
403 triggered_fetch_day = CASE
404 WHEN LEAST(triggered_fetch_at, follow_up_fetch_at) <= now() THEN (now() AT TIME ZONE 'UTC')::date
405 ELSE triggered_fetch_day
406 END,
407 triggered_fetch_at = CASE
408 WHEN triggered_fetch_at <= now() THEN NULL
409 ELSE triggered_fetch_at
410 END,
411 follow_up_fetch_at = CASE
412 WHEN follow_up_fetch_at <= now() THEN NULL
413 ELSE follow_up_fetch_at
414 END
415WHERE course_code = $1
416 "#,
417 course_code,
418 listed_person_count,
419 duration_ms,
420 )
421 .execute(conn)
422 .await?;
423 Ok(())
424}
425
426pub async fn mark_window_closed(conn: &mut PgConnection, course_code: &str) -> ModelResult<()> {
429 sqlx::query!(
430 r#"
431UPDATE credit_registration_roster_schedules
432SET last_fetched_at = now(),
433 last_listed_person_count = 0,
434 is_fetched_alone = FALSE,
435 consecutive_failures = 0,
436 retry_not_before = NULL,
437 last_error = NULL,
438 window_closed_at = now(),
439 triggered_fetch_at = NULL,
440 follow_up_fetch_at = NULL
441WHERE course_code = $1
442 "#,
443 course_code,
444 )
445 .execute(conn)
446 .await?;
447 Ok(())
448}
449
450pub async fn mark_batch_failed(
453 conn: &mut PgConnection,
454 course_codes: &[String],
455 error: CreditRegistrationErrorCode,
456) -> ModelResult<()> {
457 sqlx::query!(
458 r#"
459UPDATE credit_registration_roster_schedules
460SET is_fetched_alone = TRUE,
461 last_error = $2
462WHERE course_code = ANY($1::text [])
463 "#,
464 course_codes,
465 error as CreditRegistrationErrorCode,
466 )
467 .execute(conn)
468 .await?;
469 Ok(())
470}
471
472pub async fn mark_alone_failed(
475 conn: &mut PgConnection,
476 course_code: &str,
477 error: CreditRegistrationErrorCode,
478) -> ModelResult<()> {
479 sqlx::query!(
480 r#"
481UPDATE credit_registration_roster_schedules
482SET is_fetched_alone = TRUE,
483 consecutive_failures = consecutive_failures + 1,
484 last_error = $2,
485 retry_not_before = now() + (
486 ($3::bigint [])[LEAST(consecutive_failures + 1, CARDINALITY($3::bigint []))] * INTERVAL '1 second'
487 )
488WHERE course_code = $1
489 "#,
490 course_code,
491 error as CreditRegistrationErrorCode,
492 &ALONE_FAILURE_BACKOFF_SECS[..],
493 )
494 .execute(conn)
495 .await?;
496 Ok(())
497}
498
499pub async fn book_triggered_fetch(
504 conn: &mut PgConnection,
505 course_code: &str,
506 books_follow_up: bool,
507) -> ModelResult<bool> {
508 sqlx::query!(
509 r#"
510INSERT INTO credit_registration_roster_schedules (course_code)
511VALUES ($1) ON CONFLICT (course_code) DO NOTHING
512 "#,
513 course_code,
514 )
515 .execute(&mut *conn)
516 .await?;
517 let booked = sqlx::query!(
518 r#"
519UPDATE credit_registration_roster_schedules
520SET triggered_fetch_at = LEAST(
521 COALESCE(triggered_fetch_at, 'infinity'::timestamptz),
522 GREATEST(
523 now(),
524 COALESCE(last_fetched_at, '-infinity'::timestamptz) + ($3::bigint * INTERVAL '1 second')
525 )
526 ),
527 follow_up_fetch_at = CASE
528 WHEN $2
529 AND follow_up_fetch_at IS NULL THEN now() + ($4::bigint * INTERVAL '1 second')
530 ELSE follow_up_fetch_at
531 END
532WHERE course_code = $1
533 AND (
534 triggered_fetch_day IS DISTINCT FROM (now() AT TIME ZONE 'UTC')::date
535 OR triggered_fetch_count < $5
536 )
537 "#,
538 course_code,
539 books_follow_up,
540 TRIGGER_MIN_GAP_SECS,
541 VISIT_FOLLOW_UP_SECS,
542 MAX_TRIGGERED_FETCHES_PER_DAY,
543 )
544 .execute(conn)
545 .await?;
546 Ok(booked.rows_affected() > 0)
547}
548
549#[derive(Debug, Clone, PartialEq)]
551pub struct FailingRosterCode {
552 pub course_code: String,
553 pub consecutive_failures: i32,
554 pub last_error: Option<CreditRegistrationErrorCode>,
555 pub last_attempted_at: Option<DateTime<Utc>>,
556}
557
558pub async fn get_failing_codes(conn: &mut PgConnection) -> ModelResult<Vec<FailingRosterCode>> {
560 let res = sqlx::query_as!(
561 FailingRosterCode,
562 r#"
563SELECT course_code,
564 consecutive_failures,
565 last_error,
566 last_attempted_at
567FROM credit_registration_roster_schedules
568WHERE consecutive_failures > 0
569ORDER BY consecutive_failures DESC,
570 course_code
571 "#,
572 )
573 .fetch_all(conn)
574 .await?;
575 Ok(res)
576}
577
578pub mod testing {
580 use super::DORMANT_INTERVAL_SECS;
581 use crate::prelude::*;
582
583 pub async fn make_listings_due_for_testing(
587 conn: &mut PgConnection,
588 course_id: Uuid,
589 ) -> ModelResult<u64> {
590 let res = sqlx::query!(
591 r#"
592UPDATE credit_registration_roster_schedules
593SET last_fetched_at = last_fetched_at - ($2::bigint * INTERVAL '1 second'),
594 retry_not_before = NULL
595WHERE course_code IN (
596 SELECT TRIM(cm.uh_course_code)
597 FROM credit_registration_active_course_modules acm
598 JOIN course_modules cm ON cm.id = acm.course_module_id AND cm.deleted_at IS NULL
599 WHERE acm.course_id = $1
600 )
601 "#,
602 course_id,
603 DORMANT_INTERVAL_SECS,
604 )
605 .execute(conn)
606 .await?;
607 Ok(res.rows_affected())
608 }
609}