headless_lms_credit_registration/use_cases/verify/
poll.rs1use chrono::{DateTime, Utc};
4use headless_lms_models::credit_registrations::{
5 AdminAttention, VerifyFlow, mark_partially_registered, reset_for_resubmission,
6 set_needs_admin_attention,
7};
8use headless_lms_models::library::credit_registration::outcomes::{
9 Outcome, verify_inconclusive_outcome,
10};
11use sqlx::{Connection, PgConnection};
12
13use super::decide::{
14 PollAnswer, decide_poll, not_registered_decision, partially_registered_decision,
15};
16use super::lease::{Leased, VerifyAttempt, attempt_facts, claim_and_lease};
17use crate::error::CreditRegistrationResult;
18use crate::registry::{AttainmentId, ExchangeAudit, VerificationAnswer, VerificationRequest};
19use crate::use_cases::batch_flow::{BatchFlowContext, Prepared, RegistryBatchFlow};
20use crate::workflow::{
21 Applied, RefusalPolicy, write_decision, write_decision_committing_if_written,
22};
23
24const STUCK_WITHOUT_ATTAINMENT_ID_MESSAGE: &str =
25 "Credit registration is awaiting verification with no submitted attainment id";
26
27pub(super) struct VerifyPoll;
29
30impl RegistryBatchFlow for VerifyPoll {
31 type Extra = VerifyAttempt;
32 type Request = VerificationRequest;
33
34 const ALL_UNAVAILABLE_ERROR: &'static str = "Every verify poll came back unavailable.";
35 const REFUSAL: RefusalPolicy<VerifyAttempt> = RefusalPolicy::KeepWaiting {
39 outcome: still_polling,
40 message: "Could not verify this submission this time.",
41 };
42
43 async fn claim(
46 ctx: &BatchFlowContext<'_>,
47 conn: &mut PgConnection,
48 limit: usize,
49 ) -> CreditRegistrationResult<Prepared<VerifyAttempt, VerificationRequest>> {
50 let mut prepared = Prepared::new();
51 for poll in claim_and_lease(ctx, conn, VerifyFlow::Poll, limit).await? {
52 let row = poll.claim.registration();
53 let Some(submitted_attainment_id) = row.submitted_attainment_id.clone() else {
54 error!(
57 credit_registration_id = %row.id,
58 "Credit registration is awaiting verification with no submitted attainment id; stuck"
59 );
60 ctx.errors
61 .report(
62 STUCK_WITHOUT_ATTAINMENT_ID_MESSAGE,
63 None,
64 serde_json::json!({ "credit_registration_id": row.id }),
65 )
66 .await;
67 if !row.needs_admin_attention {
68 set_needs_admin_attention(conn, row.id, AdminAttention::Raise).await?;
69 }
70 prepared.record_failed();
71 continue;
72 };
73 let request = VerificationRequest {
74 submitted_attainment_id: AttainmentId::new(submitted_attainment_id),
75 };
76 prepared.send(poll, request);
77 }
78 Ok(prepared)
79 }
80
81 async fn apply_answer(
82 conn: &mut PgConnection,
83 poll: &Leased,
84 answer: Option<&VerificationAnswer>,
85 audit: &ExchangeAudit,
86 ) -> CreditRegistrationResult<Applied> {
87 let claim = &poll.claim;
88 let facts = attempt_facts(poll, Utc::now());
89 let row_error = answer.and_then(|answer| answer.error_message.as_deref());
90 match decide_poll(claim.registration().state, answer, &facts) {
91 PollAnswer::Decided(decision) => {
92 write_decision(conn, claim, (*decision).with_row_error(row_error), audit).await
93 }
94 PollAnswer::PartiallyRegistered => {
97 let partially_registered_at = mark_partially_registered(conn, claim.id()).await?;
98 let decision = partially_registered_decision(&facts, partially_registered_at)
99 .with_row_error(row_error);
100 write_decision(conn, claim, decision, audit).await
101 }
102 PollAnswer::NotRegistered => {
104 let mut tx = conn.begin().await?;
105 let reimport_count = reset_for_resubmission(&mut tx, claim.id()).await?;
106 let decision =
107 not_registered_decision(&facts, reimport_count).with_row_error(row_error);
108 write_decision_committing_if_written(tx, claim, decision, audit).await
109 }
110 }
111 }
112}
113
114fn still_polling(poll: &Leased, now: DateTime<Utc>) -> Outcome {
115 verify_inconclusive_outcome(poll.claim.registration().state, &attempt_facts(poll, now))
116}