Skip to main content

headless_lms_models/credit_registrations/
claims.rs

1//! Claiming due rows for a worker phase.
2//!
3//! Every claim takes up to `limit` due rows in its states, the longest due first. The row locks
4//! live until the caller's transaction ends, so callers must pass a transaction. Rows on a paused
5//! course module, or on one whose credit registration has been switched off, are never claimed:
6//! enforced here so no phase can forget it. Both freeze a row where it stands rather than
7//! cancelling it, so switching the module back on resumes the rows that were already in flight.
8//!
9//! An unscoped claim (the live background worker) also skips a row whose user (and, if the hold
10//! names one, course) has a live row in `credit_registration_test_exclusive_holds`. A scoped claim
11//! always ignores holds, so a spec driving its own rows through explicit ticks is unaffected
12//! either way.
13
14use crate::prelude::*;
15use chrono::TimeDelta;
16use secrecy::ExposeSecret;
17
18use super::registration::CreditRegistration;
19use super::state::CreditRegistrationState;
20
21/// Which rows a phase iteration may touch. Empty means every row, which is what production runs; a
22/// narrowed scope lets a test drive the pipeline for its own course on a shared database.
23#[derive(Debug, Clone, Default, PartialEq)]
24pub struct RegistrationScope {
25    pub course_id: Option<Uuid>,
26    pub user_id: Option<Uuid>,
27    /// The precision escape hatch, for a caller that already knows its ledger rows.
28    pub credit_registration_ids: Vec<Uuid>,
29}
30
31impl RegistrationScope {
32    pub fn is_unscoped(&self) -> bool {
33        self.course_id.is_none()
34            && self.user_id.is_none()
35            && self.credit_registration_ids.is_empty()
36    }
37
38    pub fn for_course(course_id: Uuid) -> Self {
39        Self {
40            course_id: Some(course_id),
41            ..Self::default()
42        }
43    }
44}
45
46/// The states resolve-enrolments looks up from: rows on their way to a first resolve, and rows
47/// parked in `no_usable_enrolment` whose next enrolment check is due, which are checked where they
48/// stand.
49const LOOKUP_STATES: [CreditRegistrationState; 2] = [
50    CreditRegistrationState::ReadyToSubmit,
51    CreditRegistrationState::NoUsableEnrolment,
52];
53
54/// Claims, for the person lookup that precedes resolve-enrolments, the rows
55/// [`claim_due_for_resolve`] would take, the later [`EnrolmentCheckGroup`](crate::library::credit_registration::enrolment_check_schedule::EnrolmentCheckGroup) first.
56pub async fn claim_due_for_person_lookup(
57    conn: &mut PgConnection,
58    scope: &RegistrationScope,
59    limit: i64,
60) -> ModelResult<Vec<CreditRegistration>> {
61    claim(conn, &LOOKUP_STATES, scope, limit, ClaimKind::PersonLookup).await
62}
63
64/// Claims, for resolve-enrolments, `ready_to_submit` rows and parked rows due an enrolment check,
65/// the later [`EnrolmentCheckGroup`](crate::library::credit_registration::enrolment_check_schedule::EnrolmentCheckGroup) first, minus any whose student already has another live row
66/// for the module somewhere between resolving and a known outcome. First pulls slow checks due soon
67/// into the batch; see
68/// [`crate::library::credit_registration::enrolment_checks::pull_forward_batched_checks`].
69///
70/// Only one completion per student and module goes past resolve at a time, so each is weighed
71/// against the outcome of the one before it rather than racing it to the registry; see
72/// [`lock_live_successes_for_same_module`](super::lock_live_successes_for_same_module). Of two such
73/// rows claimed together only the first comes back; the other stays claimable where it is. A parked
74/// row sent for a check must be marked with [`claim_enrolment_checks`] before the claim's
75/// transaction commits.
76pub async fn claim_due_for_resolve(
77    conn: &mut PgConnection,
78    scope: &RegistrationScope,
79    limit: i64,
80) -> ModelResult<Vec<CreditRegistration>> {
81    crate::library::credit_registration::enrolment_checks::pull_forward_batched_checks(conn, scope)
82        .await?;
83    let claimed = claim(conn, &LOOKUP_STATES, scope, limit, ClaimKind::Resolve).await?;
84    Ok(first_per(
85        claimed,
86        |row| (row.user_id, row.course_module_id),
87        "another attempt for the same student and module",
88    ))
89}
90
91/// Keeps parked rows out of every claim while a lookup for them is out, as `resolving_enrolment`
92/// does for a row on its first resolve. The answer's [`transition`](super::transition::transition) ends the
93/// claim; one a worker died holding expires after
94/// [`RESOLVING_RECOVERY_GRACE`](crate::library::credit_registration::backoff::RESOLVING_RECOVERY_GRACE).
95/// A no-op for a row in any other state.
96pub async fn claim_enrolment_checks(conn: &mut PgConnection, ids: &[Uuid]) -> ModelResult<()> {
97    use crate::library::credit_registration::backoff::RESOLVING_RECOVERY_GRACE;
98    if ids.is_empty() {
99        return Ok(());
100    }
101    sqlx::query!(
102        r#"
103UPDATE credit_registrations
104SET enrolment_check_claimed_until = now() + $2::interval
105WHERE id = ANY($1)
106  AND state = 'no_usable_enrolment'
107        "#,
108        ids,
109        RESOLVING_RECOVERY_GRACE as TimeDelta,
110    )
111    .execute(conn)
112    .await?;
113    Ok(())
114}
115
116/// Claims, for import, `checking_enrolment` rows, minus any whose student and course code already
117/// have a submission in flight, which Suotar's hour-old copy of Sisu would not stop from
118/// registering twice, and any whose person and module slot in
119/// `uq_credit_registrations_person_module` is taken. Of two rows for the same student and course
120/// code claimed together only the first comes back; the other stays claimable where it is.
121pub async fn claim_due_for_import(
122    conn: &mut PgConnection,
123    scope: &RegistrationScope,
124    limit: i64,
125) -> ModelResult<Vec<CreditRegistration>> {
126    let claimed = claim(
127        conn,
128        &[CreditRegistrationState::CheckingEnrolment],
129        scope,
130        limit,
131        ClaimKind::Import,
132    )
133    .await?;
134    Ok(first_per(
135        claimed,
136        |row| {
137            (
138                row.student_number
139                    .as_ref()
140                    .map(|number| number.expose_secret().to_string()),
141                row.uh_course_code.clone(),
142            )
143        },
144        "another attempt for the same student and course code",
145    ))
146}
147
148/// The first of the claimed rows per `key`. The claim query holds back a row whose twin is already
149/// in flight, but cannot see a twin claimed alongside it; the rows dropped here keep their lock
150/// until the claim's transaction ends, and are claimable where they are after it.
151fn first_per<K: Eq + std::hash::Hash>(
152    claimed: Vec<CreditRegistration>,
153    key: impl Fn(&CreditRegistration) -> K,
154    twin: &str,
155) -> Vec<CreditRegistration> {
156    let mut seen = std::collections::HashSet::new();
157    claimed
158        .into_iter()
159        .filter(|row| {
160            let is_first = seen.insert(key(row));
161            if !is_first {
162                debug!(
163                    credit_registration_id = %row.id,
164                    "Leaving row claimable: {twin} is already in this claim"
165                );
166            }
167            is_first
168        })
169        .collect()
170}
171
172/// The states verify polls from. Withdrawal moves a row out of all of them, which is what stops the
173/// polling without any query having to know about withdrawal.
174const VERIFY_STATES: [CreditRegistrationState; 3] = [
175    CreditRegistrationState::AwaitingVerification,
176    CreditRegistrationState::PartiallyRegistered,
177    CreditRegistrationState::SubmissionUncertain,
178];
179
180/// Which of verify's flows a claim is for.
181#[derive(Debug, Clone, Copy, PartialEq, Eq)]
182pub enum VerifyFlow {
183    /// Rows with a submitted attainment to poll by, and `awaiting_verification` rows stuck without
184    /// one.
185    Poll,
186    /// `submission_uncertain` rows with no submitted attainment id, looked up through
187    /// resolve-enrolments instead.
188    UncertainRecovery,
189}
190
191/// Claims rows for one of verify's flows. Each flow is claimed on its own, so that rows one flow
192/// has no allowance to send cannot fill the other's claim.
193pub async fn claim_due_for_verify(
194    conn: &mut PgConnection,
195    flow: VerifyFlow,
196    scope: &RegistrationScope,
197    limit: i64,
198) -> ModelResult<Vec<CreditRegistration>> {
199    claim(conn, &VERIFY_STATES, scope, limit, ClaimKind::Verify(flow)).await
200}
201
202/// Which caller a claim is for, which decides the rows that hold a row back and the order.
203#[derive(Debug, Clone, Copy, PartialEq, Eq)]
204enum ClaimKind {
205    /// See [`claim_due_for_verify`].
206    Verify(VerifyFlow),
207    /// See [`claim_due_for_person_lookup`].
208    PersonLookup,
209    /// See [`claim_due_for_resolve`].
210    Resolve,
211    /// See [`claim_due_for_import`].
212    Import,
213}
214
215/// Shares its eligibility filters with
216/// [`count_due_enrolment_checks`](super::count_due_enrolment_checks) and
217/// [`pull_forward_batched_checks`](crate::library::credit_registration::enrolment_checks::pull_forward_batched_checks);
218/// change all three together.
219async fn claim(
220    conn: &mut PgConnection,
221    states: &[CreditRegistrationState],
222    scope: &RegistrationScope,
223    limit: i64,
224    kind: ClaimKind,
225) -> ModelResult<Vec<CreditRegistration>> {
226    let ignores_test_holds = !scope.is_unscoped();
227    let excludes_twins_in_flight = kind == ClaimKind::Import;
228    let waits_behind_other_attempts = kind == ClaimKind::Resolve;
229    let orders_by_check_group = matches!(kind, ClaimKind::PersonLookup | ClaimKind::Resolve);
230    // `Some(true)` claims only uncertain rows without a submitted attainment, `Some(false)` every
231    // other verify row.
232    let only_uncertain_without_attainment_id = match kind {
233        ClaimKind::Verify(flow) => Some(flow == VerifyFlow::UncertainRecovery),
234        _ => None,
235    };
236    let res = sqlx::query_as!(
237        CreditRegistration,
238        r#"
239WITH due AS (
240  SELECT cr.id
241  FROM credit_registrations cr
242    JOIN credit_registration_active_course_modules acm ON acm.course_module_id = cr.course_module_id
243    JOIN course_module_completions cmc ON cmc.id = cr.course_module_completion_id
244  WHERE cr.deleted_at IS NULL
245    AND cr.superseded_by_id IS NULL
246    -- A completion opted out by hand is the pull path's again; only a row already sent carries on.
247    AND (
248      cmc.register_credits_via_suotar
249      OR cr.submitted_at IS NOT NULL
250    )
251    AND cr.state = ANY($1::credit_registration_state [])
252    AND cr.next_attempt_at <= now()
253    AND (
254      cr.enrolment_check_claimed_until IS NULL
255      OR cr.enrolment_check_claimed_until <= now()
256    )
257    AND ($3::uuid IS NULL OR cr.course_id = $3)
258    AND ($4::uuid IS NULL OR cr.user_id = $4)
259    AND (
260      cardinality($5::uuid []) = 0
261      OR cr.id = ANY($5::uuid [])
262    )
263    AND (
264      $12::boolean IS NULL
265      OR (
266        cr.state = 'submission_uncertain'
267        AND cr.submitted_attainment_id IS NULL
268      ) = $12
269    )
270    AND (
271      $6::boolean
272      OR NOT EXISTS (
273        SELECT 1
274        FROM credit_registration_test_exclusive_holds h
275        WHERE h.user_id = cr.user_id
276          AND (
277            h.course_id IS NULL
278            OR h.course_id = cr.course_id
279          )
280          AND h.held_until > now()
281      )
282    )
283    AND (
284      NOT $7::boolean
285      OR (
286        NOT EXISTS (
287          SELECT 1
288          FROM credit_registrations twin
289          WHERE twin.student_number = cr.student_number
290            AND twin.uh_course_code = cr.uh_course_code
291            AND twin.id <> cr.id
292            AND twin.deleted_at IS NULL
293            AND twin.state = ANY($10::credit_registration_state [])
294        )
295        -- Mirrors uq_credit_registrations_person_module, which moving to submitting would violate.
296        AND NOT EXISTS (
297          SELECT 1
298          FROM credit_registrations holder
299          WHERE holder.sisu_person_id = cr.sisu_person_id
300            AND holder.course_module_id = cr.course_module_id
301            AND holder.id <> cr.id
302            AND holder.deleted_at IS NULL
303            AND holder.superseded_by_id IS NULL
304            AND holder.pending_superseded_by_id IS NULL
305            AND (
306              holder.state = ANY($10::credit_registration_state [])
307              OR holder.state = ANY($11::credit_registration_state [])
308            )
309        )
310      )
311    )
312    AND (
313      NOT $8::boolean
314      OR NOT EXISTS (
315        SELECT 1
316        FROM credit_registrations ahead
317        WHERE ahead.user_id = cr.user_id
318          AND ahead.course_module_id = cr.course_module_id
319          AND ahead.id <> cr.id
320          AND ahead.deleted_at IS NULL
321          AND ahead.superseded_by_id IS NULL
322          -- failed_retryable because its backoff may resume it at checking_enrolment or later.
323          AND (
324            ahead.state IN (
325              'resolving_enrolment',
326              'checking_enrolment',
327              'failed_retryable'
328            )
329            OR ahead.state = ANY($10::credit_registration_state [])
330            OR ahead.enrolment_check_claimed_until > now()
331          )
332      )
333    )
334  ORDER BY CASE
335      WHEN $9 THEN cr.enrolment_check_group
336    END DESC NULLS LAST,
337    cr.next_attempt_at
338  FOR UPDATE OF cr SKIP LOCKED
339  LIMIT $2
340)
341UPDATE credit_registrations cr
342SET last_attempt_at = now()
343FROM due
344WHERE cr.id = due.id
345RETURNING cr.*
346        "#,
347        states as &[CreditRegistrationState],
348        limit,
349        scope.course_id,
350        scope.user_id,
351        &scope.credit_registration_ids,
352        ignores_test_holds,
353        excludes_twins_in_flight,
354        waits_behind_other_attempts,
355        orders_by_check_group,
356        &CreditRegistrationState::IN_FLIGHT_STATES as &[CreditRegistrationState],
357        &CreditRegistrationState::SUCCESS_STATES as &[CreditRegistrationState],
358        only_uncertain_without_attainment_id,
359    )
360    .fetch_all(conn)
361    .await?;
362    Ok(res)
363}