headless_lms_credit_registration/use_cases/import/
mod.rs1mod claim;
8mod decide;
9
10use headless_lms_models::credit_registrations::restamp_submitting;
11use headless_lms_models::library::credit_registration::outcomes::released_unsent_split_half;
12use headless_lms_utils::prelude::Utc;
13use sqlx::PgConnection;
14
15use crate::error::CreditRegistrationResult;
16use crate::registry::{AttainmentSubmission, ExchangeAudit, ImportAnswer, StudyRegistry};
17use crate::use_cases::batch_flow::{
18 BatchFlowContext, Prepared, RegistryBatchFlow, run_registry_batch_flow,
19};
20use crate::workflow::{
21 Applied, Claimed, Counts, RefusalPolicy, write_decision, write_unasked_move,
22};
23
24use claim::claim_import_candidates;
25use decide::decide_import_answer;
26
27pub(crate) async fn run<R: StudyRegistry>(
28 ctx: &BatchFlowContext<'_>,
29 registry: &mut R,
30) -> CreditRegistrationResult<Counts> {
31 run_registry_batch_flow::<Import, _>(ctx, registry).await
32}
33
34struct Import;
35
36impl RegistryBatchFlow for Import {
37 type Extra = ();
38 type Request = AttainmentSubmission;
39
40 const ALL_UNAVAILABLE_ERROR: &'static str =
41 "Every item of the batch timed out in Sisu or came back unavailable.";
42 const REFUSAL: RefusalPolicy<()> = RefusalPolicy::RequestLevel;
43
44 async fn claim(
45 ctx: &BatchFlowContext<'_>,
46 conn: &mut PgConnection,
47 limit: usize,
48 ) -> CreditRegistrationResult<Prepared<(), AttainmentSubmission>> {
49 claim_import_candidates(ctx, conn, limit).await
50 }
51
52 async fn apply_answer(
53 conn: &mut PgConnection,
54 row: &Claimed<()>,
55 answer: Option<&ImportAnswer>,
56 audit: &ExchangeAudit,
57 ) -> CreditRegistrationResult<Applied> {
58 let claim = &row.claim;
59 let decision = decide_import_answer(claim.registration(), answer, &claim.facts(Utc::now()));
60 write_decision(conn, claim, decision, audit).await
61 }
62
63 async fn keep_in_flight(
64 conn: &mut PgConnection,
65 rows: &[&Claimed<()>],
66 ) -> CreditRegistrationResult<()> {
67 restamp_submitting(
68 conn,
69 &rows.iter().map(|row| row.claim.id()).collect::<Vec<_>>(),
70 )
71 .await?;
72 Ok(())
73 }
74
75 async fn release_unsent(
76 conn: &mut PgConnection,
77 rows: &[&Claimed<()>],
78 ) -> CreditRegistrationResult<()> {
79 for row in rows {
80 write_unasked_move(conn, &row.claim, released_unsent_split_half()).await?;
81 }
82 Ok(())
83 }
84}