headless_lms_credit_registration/use_cases/verify/
lease.rs1use 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
16pub(super) type Leased = Claimed<VerifyAttempt>;
18
19pub(super) struct VerifyAttempt(i32);
22
23pub(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
32pub(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 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}