Skip to main content

headless_lms_models/credit_registrations/
row_writes.rs

1//! Writes to a row's fields other than its state, which only [`transition`](super::transition::transition)
2//! writes: the frozen payload, what the study registry answered, the schedule, the counters and the flags.
3
4use super::transition::AdminAttention;
5use crate::library::credit_registration::enrolment_check_schedule::EnrolmentCheckSource;
6use crate::prelude::*;
7use headless_lms_utils::secret_string::expose_option;
8use secrecy::ExposeSecret;
9use std::collections::HashMap;
10
11/// Frozen copy of what we are about to submit. Written once, before the row leaves
12/// `checking_enrolment`: a later regrade must not alter a submitted row.
13#[derive(Debug, Clone)]
14pub struct PayloadSnapshot {
15    pub student_number: DbSecret,
16    pub sisu_person_id: Option<DbSecret>,
17    pub uh_course_code: String,
18    pub selected_enrolment_id: Option<String>,
19    pub selected_enrolment_kind: Option<String>,
20    pub selected_enrolment_realisation_id: Option<String>,
21    /// Localized `{fi, sv, en}` realisation name, as Suotar reported it.
22    pub selected_enrolment_realisation_name: Option<serde_json::Value>,
23    pub attained_at: DateTime<Utc>,
24    pub attainment_language: String,
25    pub grade_scale_id: String,
26    pub grade_id: String,
27    pub credits: f32,
28}
29
30pub async fn set_payload_snapshot(
31    conn: &mut PgConnection,
32    id: Uuid,
33    snapshot: &PayloadSnapshot,
34) -> ModelResult<()> {
35    sqlx::query!(
36        r#"
37UPDATE credit_registrations
38SET student_number = $2,
39  sisu_person_id = $3,
40  uh_course_code = $4,
41  selected_enrolment_id = $5,
42  selected_enrolment_kind = $6,
43  selected_enrolment_realisation_id = $7,
44  attained_at = $8,
45  attainment_language = $9,
46  grade_scale_id = $10,
47  grade_id = $11,
48  credits = $12,
49  selected_enrolment_realisation_name = $13
50WHERE id = $1
51  AND deleted_at IS NULL
52        "#,
53        id,
54        snapshot.student_number.expose_secret(),
55        expose_option(&snapshot.sisu_person_id),
56        snapshot.uh_course_code,
57        snapshot.selected_enrolment_id,
58        snapshot.selected_enrolment_kind,
59        snapshot.selected_enrolment_realisation_id,
60        snapshot.attained_at,
61        snapshot.attainment_language,
62        snapshot.grade_scale_id,
63        snapshot.grade_id,
64        snapshot.credits,
65        snapshot.selected_enrolment_realisation_name,
66    )
67    .execute(conn)
68    .await?;
69    Ok(())
70}
71
72pub async fn set_submitted_attainment(
73    conn: &mut PgConnection,
74    id: Uuid,
75    submitted_attainment_id: &str,
76    submitted_attainment_type: Option<&str>,
77) -> ModelResult<()> {
78    sqlx::query!(
79        r#"
80UPDATE credit_registrations
81SET submitted_attainment_id = $2,
82  submitted_attainment_type = $3
83WHERE id = $1
84  AND deleted_at IS NULL
85        "#,
86        id,
87        submitted_attainment_id,
88        submitted_attainment_type,
89    )
90    .execute(conn)
91    .await?;
92    Ok(())
93}
94
95/// Notes that verify saw only the assessment item attainment, keeping the first sighting, and
96/// returns when that was.
97pub async fn mark_partially_registered(
98    conn: &mut PgConnection,
99    id: Uuid,
100) -> ModelResult<DateTime<Utc>> {
101    let partially_registered_at = sqlx::query_scalar!(
102        r#"
103UPDATE credit_registrations
104SET partially_registered_at = COALESCE(partially_registered_at, now())
105WHERE id = $1
106  AND deleted_at IS NULL
107RETURNING partially_registered_at AS "partially_registered_at!"
108        "#,
109        id,
110    )
111    .fetch_one(conn)
112    .await?;
113    Ok(partially_registered_at)
114}
115
116/// Restamps `submitted_at` on rows still `submitting`, for an import that sends them again after
117/// splitting a refused batch: the precondition sweep times a lost submission from this stamp, and
118/// must not condemn a row still waiting its turn in the same iteration.
119pub async fn restamp_submitting(conn: &mut PgConnection, ids: &[Uuid]) -> ModelResult<()> {
120    sqlx::query!(
121        r#"
122UPDATE credit_registrations
123SET submitted_at = now()
124WHERE id = ANY($1)
125  AND state = 'submitting'
126  AND deleted_at IS NULL
127        "#,
128        ids,
129    )
130    .execute(conn)
131    .await?;
132    Ok(())
133}
134
135/// Restarts the recovery grace of rows waiting out an enrolment lookup in `resolving_enrolment`,
136/// which runs from `state_entered_at`, for a split batch whose later halves are still to be sent.
137pub async fn restamp_resolving_enrolment(conn: &mut PgConnection, ids: &[Uuid]) -> ModelResult<()> {
138    if ids.is_empty() {
139        return Ok(());
140    }
141    sqlx::query!(
142        r#"
143UPDATE credit_registrations
144SET state_entered_at = now()
145WHERE id = ANY($1)
146  AND state = 'resolving_enrolment'
147  AND deleted_at IS NULL
148        "#,
149        ids,
150    )
151    .execute(conn)
152    .await?;
153    Ok(())
154}
155
156/// Records Suotar's `retryAfter` for the pending submission, before which a resubmission may
157/// register the credits twice. See
158/// [`ResubmissionFacts::resubmission_refusal`](super::ResubmissionFacts::resubmission_refusal).
159pub async fn set_resubmit_not_before(
160    conn: &mut PgConnection,
161    id: Uuid,
162    resubmit_not_before: DateTime<Utc>,
163) -> ModelResult<()> {
164    sqlx::query!(
165        r#"
166UPDATE credit_registrations
167SET resubmit_not_before = $2
168WHERE id = $1
169  AND deleted_at IS NULL
170        "#,
171        id,
172        resubmit_not_before,
173    )
174    .execute(conn)
175    .await?;
176    Ok(())
177}
178
179/// Forgets a submission Suotar says never landed, so the row resolves its enrolment and imports
180/// again from scratch, and returns how many times that has now happened.
181///
182/// Also restarts the retry window and the verify count: the resend is new work, and the old
183/// submission's history would expire it or flag it at once.
184pub async fn reset_for_resubmission(conn: &mut PgConnection, id: Uuid) -> ModelResult<i32> {
185    let reimport_count = sqlx::query_scalar!(
186        r#"
187UPDATE credit_registrations
188SET submitted_attainment_id = NULL,
189  submitted_attainment_type = NULL,
190  partially_registered_at = NULL,
191  resubmit_not_before = NULL,
192  selected_enrolment_id = NULL,
193  grade_id = NULL,
194  first_failed_at = NULL,
195  verify_attempt_count = 0,
196  not_registered_reimport_count = not_registered_reimport_count + 1
197WHERE id = $1
198  AND deleted_at IS NULL
199RETURNING not_registered_reimport_count
200        "#,
201        id,
202    )
203    .fetch_one(conn)
204    .await?;
205    Ok(reimport_count)
206}
207
208/// Records the attainment the study registry holds, unless another live row already claims it.
209///
210/// Two rows may legitimately be told about one attainment — a grade improvement Sisu declines names
211/// the attainment the first attempt registered — so this returns `false` instead of failing.
212pub async fn set_sisu_attainment_if_unclaimed(
213    conn: &mut PgConnection,
214    id: Uuid,
215    sisu_attainment_id: &str,
216    sisu_attainment_type: Option<&str>,
217) -> ModelResult<bool> {
218    let updated = sqlx::query_scalar!(
219        r#"
220UPDATE credit_registrations
221SET sisu_attainment_id = $2,
222  sisu_attainment_type = $3
223WHERE id = $1
224  AND deleted_at IS NULL
225  AND NOT EXISTS (
226    SELECT 1
227    FROM credit_registrations other
228    WHERE other.sisu_attainment_id = $2
229      AND other.deleted_at IS NULL
230      AND other.id <> $1
231  )
232RETURNING id
233        "#,
234        id,
235        sisu_attainment_id,
236        sisu_attainment_type,
237    )
238    .fetch_optional(conn)
239    .await;
240    match updated {
241        Ok(updated) => Ok(updated.is_some()),
242        // The NOT EXISTS guard above isn't atomic against a concurrent caller claiming the
243        // same sisu_attainment_id for a different row; the loser hits this unique index instead.
244        Err(err) => {
245            let err: ModelError = err.into();
246            match err.error_type() {
247                ModelErrorType::DatabaseConstraint { constraint, .. }
248                    if constraint == "uq_credit_registrations_sisu_attainment" =>
249                {
250                    Ok(false)
251                }
252                _ => Err(err),
253            }
254        }
255    }
256}
257
258/// Defers when the pipeline may next claim this row; the delay is the caller's policy.
259///
260/// Deliberately leaves `first_failed_at` alone: a state that waits on a human defers too, and
261/// anchoring the retry window here would expire it.
262pub async fn schedule_next_attempt(
263    conn: &mut PgConnection,
264    id: Uuid,
265    next_attempt_at: DateTime<Utc>,
266) -> ModelResult<()> {
267    sqlx::query!(
268        r#"
269UPDATE credit_registrations
270SET next_attempt_at = $2
271WHERE id = $1
272  AND deleted_at IS NULL
273        "#,
274        id,
275        next_attempt_at,
276    )
277    .execute(conn)
278    .await?;
279    Ok(())
280}
281
282/// Makes rows claimable again now, whatever backoff parked them. A row waiting for an enrolment
283/// check has its check marked as `enrolment_check_source`: who asked for it.
284///
285/// Uses the database clock: an app-clock value sampled after `BEGIN` is still in the future when
286/// the same transaction compares it against `now()`.
287pub async fn make_due_now_batch(
288    conn: &mut PgConnection,
289    ids: &[Uuid],
290    enrolment_check_source: EnrolmentCheckSource,
291) -> ModelResult<()> {
292    sqlx::query!(
293        r#"
294UPDATE credit_registrations
295SET next_attempt_at = now(),
296  enrolment_check_source = CASE
297    WHEN state = 'no_usable_enrolment' THEN $2
298    ELSE enrolment_check_source
299  END
300WHERE id = ANY($1)
301  AND next_attempt_at > now()
302  AND superseded_by_id IS NULL
303  AND deleted_at IS NULL
304        "#,
305        ids,
306        enrolment_check_source as EnrolmentCheckSource,
307    )
308    .execute(conn)
309    .await?;
310    Ok(())
311}
312
313pub async fn increment_submit_retry_count(conn: &mut PgConnection, id: Uuid) -> ModelResult<i32> {
314    let res = sqlx::query!(
315        r#"
316UPDATE credit_registrations
317SET submit_retry_count = submit_retry_count + 1
318WHERE id = $1
319  AND deleted_at IS NULL
320RETURNING submit_retry_count
321        "#,
322        id
323    )
324    .fetch_one(conn)
325    .await?;
326    Ok(res.submit_retry_count)
327}
328
329/// Counts one verify poll for every row of a batch and returns each row's new count, which sets the
330/// backoff the poll's answer is scheduled by.
331pub async fn increment_verify_attempt_counts(
332    conn: &mut PgConnection,
333    ids: &[Uuid],
334) -> ModelResult<HashMap<Uuid, i32>> {
335    let rows = sqlx::query!(
336        r#"
337UPDATE credit_registrations
338SET verify_attempt_count = verify_attempt_count + 1
339WHERE id = ANY($1)
340  AND deleted_at IS NULL
341RETURNING id,
342  verify_attempt_count
343        "#,
344        ids
345    )
346    .fetch_all(conn)
347    .await?;
348    Ok(rows
349        .into_iter()
350        .map(|row| (row.id, row.verify_attempt_count))
351        .collect())
352}
353
354/// [`schedule_next_attempt`] for a whole batch, each row with its own time.
355pub async fn schedule_next_attempts(
356    conn: &mut PgConnection,
357    scheduled: &[(Uuid, DateTime<Utc>)],
358) -> ModelResult<()> {
359    let (ids, times): (Vec<Uuid>, Vec<DateTime<Utc>>) = scheduled.iter().copied().unzip();
360    sqlx::query!(
361        r#"
362UPDATE credit_registrations cr
363SET next_attempt_at = scheduled.at
364FROM UNNEST($1::uuid [], $2::timestamptz []) AS scheduled(id, at)
365WHERE cr.id = scheduled.id
366  AND cr.deleted_at IS NULL
367        "#,
368        &ids,
369        &times,
370    )
371    .execute(conn)
372    .await?;
373    Ok(())
374}
375
376pub async fn set_needs_admin_attention(
377    conn: &mut PgConnection,
378    id: Uuid,
379    attention: AdminAttention,
380) -> ModelResult<()> {
381    sqlx::query!(
382        r#"
383UPDATE credit_registrations
384SET needs_admin_attention = $2
385WHERE id = $1
386  AND deleted_at IS NULL
387        "#,
388        id,
389        attention.is_raised(),
390    )
391    .execute(conn)
392    .await?;
393    Ok(())
394}
395
396/// Points an old attempt at the newer one that replaced it. The old row keeps its state and
397/// `terminal_at`: it really was registered.
398///
399/// For fixtures planting a finished replacement. The pipeline goes through
400/// [`mark_pending_superseded`], which supersedes the row only once the new attempt is registered.
401///
402/// `superseded_by_id` may name a row that does not exist yet, as long as it is inserted before the
403/// caller's transaction commits: the foreign key is deferred.
404pub async fn mark_superseded(
405    conn: &mut PgConnection,
406    id: Uuid,
407    superseded_by_id: Uuid,
408) -> ModelResult<()> {
409    sqlx::query!(
410        r#"
411UPDATE credit_registrations
412SET superseded_by_id = $2,
413  superseded_at = now()
414WHERE id = $1
415  AND deleted_at IS NULL
416        "#,
417        id,
418        superseded_by_id,
419    )
420    .execute(conn)
421    .await?;
422    Ok(())
423}
424
425/// Marks a registered row as being replaced by `superseded_by_id`, a later attempt with a better
426/// grade, of the same completion or another, which takes over the row's slot in
427/// `uq_credit_registrations_person_module`.
428///
429/// Not [`mark_superseded`]: the row stays the live credit until the new attempt is registered,
430/// since Sisu may accept the better grade without ever making it the course unit's attainment.
431/// [`transition`](super::transition::transition) then supersedes the row, or clears the mark if
432/// the new attempt stops short.
433pub async fn mark_pending_superseded(
434    conn: &mut PgConnection,
435    id: Uuid,
436    superseded_by_id: Uuid,
437) -> ModelResult<()> {
438    sqlx::query!(
439        r#"
440UPDATE credit_registrations
441SET pending_superseded_by_id = $2
442WHERE id = $1
443  AND deleted_at IS NULL
444        "#,
445        id,
446        superseded_by_id,
447    )
448    .execute(conn)
449    .await?;
450    Ok(())
451}
452
453/// Records that the grade-improvement scan looked at this accepted attempt against a completion in
454/// the given revision and found nothing better.
455///
456/// `completion_updated_at` must be the `updated_at` the scan actually read, not `now()`: the point
457/// is that the row stops being a candidate until the completion changes again.
458pub async fn mark_improvement_checked(
459    conn: &mut PgConnection,
460    id: Uuid,
461    completion_updated_at: DateTime<Utc>,
462) -> ModelResult<()> {
463    sqlx::query!(
464        r#"
465UPDATE credit_registrations
466SET improvement_checked_completion_updated_at = $2
467WHERE id = $1
468  AND deleted_at IS NULL
469        "#,
470        id,
471        completion_updated_at,
472    )
473    .execute(conn)
474    .await?;
475    Ok(())
476}
477
478/// Makes every due-later `failed_retryable` row due now; returns how many. The button pressed once
479/// the study registry says an outage is over.
480///
481/// Only `failed_retryable`: no other state's backoff means "waiting out an outage", and
482/// `submission_uncertain` must never be swept forward in bulk.
483pub async fn requeue_retryable_now(
484    conn: &mut PgConnection,
485    course_id: Option<Uuid>,
486    course_module_id: Option<Uuid>,
487    limit: i64,
488) -> ModelResult<i64> {
489    let count = sqlx::query_scalar!(
490        r#"
491WITH due AS (
492  SELECT id
493  FROM credit_registrations
494  WHERE state = 'failed_retryable'
495    AND next_attempt_at > now()
496    AND superseded_by_id IS NULL
497    AND deleted_at IS NULL
498    AND ($2::uuid IS NULL OR course_id = $2)
499    AND ($3::uuid IS NULL OR course_module_id = $3)
500  ORDER BY next_attempt_at
501  LIMIT $1
502),
503updated AS (
504  UPDATE credit_registrations cr
505  SET next_attempt_at = now()
506  FROM due
507  WHERE cr.id = due.id
508  RETURNING cr.id
509)
510SELECT COUNT(*) AS "count!"
511FROM updated
512        "#,
513        limit,
514        course_id,
515        course_module_id,
516    )
517    .fetch_one(conn)
518    .await?;
519    Ok(count)
520}