Skip to main content

headless_lms_credit_registration/use_cases/verify/
lease.rs

1//! Leasing the rows a verify iteration asks about, so no concurrent iteration asks about them too.
2
3use chrono::{DateTime, Utc};
4use headless_lms_models::credit_registrations::{
5    VerifyFlow, claim_due_for_verify, increment_verify_attempt_counts, schedule_next_attempts,
6};
7use headless_lms_models::library::credit_registration::outcomes::{
8    RowFacts, verify_poll_lease_until,
9};
10use sqlx::PgConnection;
11
12use crate::error::CreditRegistrationResult;
13use crate::use_cases::batch_flow::BatchFlowContext;
14use crate::workflow::{Claimed, ClaimedRegistration};
15
16/// A claimed row, leased where it stands, with the verify attempt the lease counted for it.
17pub(super) type Leased = Claimed<VerifyAttempt>;
18
19/// The verify attempt a lease counted for its row, which sets the backoff an uncertain row's answer
20/// is scheduled by. Only [`claim_and_lease`] makes one, so it is always the count the lease wrote.
21pub(super) struct VerifyAttempt(i32);
22
23/// The row's facts under the attempt the lease counted, not the count the row was claimed with, so
24/// the backoff advances once per attempt.
25pub(super) fn attempt_facts(leased: &Leased, now: DateTime<Utc>) -> RowFacts {
26    RowFacts {
27        verify_attempt_count: leased.extra.0,
28        ..leased.claim.facts(now)
29    }
30}
31
32/// Claims up to `limit` of `flow`'s due rows, counts an attempt on each and leases it until its
33/// next poll would be due, or the iteration's registry calls run out, so a concurrent iteration cannot
34/// poll the same row. Each answer overwrites its own row's schedule.
35pub(super) async fn claim_and_lease(
36    ctx: &BatchFlowContext<'_>,
37    conn: &mut PgConnection,
38    flow: VerifyFlow,
39    limit: usize,
40) -> CreditRegistrationResult<Vec<Leased>> {
41    let claimed = claim_due_for_verify(
42        conn,
43        flow,
44        ctx.scope,
45        i64::try_from(limit).unwrap_or(i64::MAX),
46    )
47    .await?;
48    let attempts = increment_verify_attempt_counts(
49        conn,
50        &claimed.iter().map(|row| row.id).collect::<Vec<_>>(),
51    )
52    .await?;
53    let now = Utc::now();
54    let scheduled: Vec<_> = claimed
55        .iter()
56        .filter(|row| attempts.contains_key(&row.id))
57        .map(|row| {
58            (
59                row.id,
60                verify_poll_lease_until(now, row.submitted_at, ctx.study_registry_wait),
61            )
62        })
63        .collect();
64    schedule_next_attempts(conn, &scheduled).await?;
65    Ok(claimed
66        .into_iter()
67        .filter_map(|row| {
68            // Not expected while the claim holds the row's lock; if it happens anyway, the row is
69            // left unpolled rather than polled under an attempt the lease never wrote.
70            let Some(attempt) = attempts.get(&row.id).copied() else {
71                warn!(
72                    credit_registration_id = %row.id,
73                    "Verify attempt was not counted for a claimed credit registration; leaving it unpolled"
74                );
75                return None;
76            };
77            Some(Claimed {
78                claim: ClaimedRegistration::left_in_place(row),
79                extra: VerifyAttempt(attempt),
80            })
81        })
82        .collect())
83}