headless_lms_server/domain/credit_registration_phases/
product_token_refresh.rs1use headless_lms_models::course_module_suotar_configurations::get_stalest_product_ids_for_enabled_modules;
7use headless_lms_models::credit_registration_events::scrub_text;
8use headless_lms_models::credit_registration_phase_state::PhaseRunOutcome;
9use headless_lms_models::open_university_product_access_tokens::{
10 NewOpenUniversityProductAccessToken, record_refresh_failure, upsert,
11};
12use headless_lms_models::secret::DbSecret;
13use headless_lms_utils::prelude::BackendError;
14use headless_lms_utils::services::suotar::{
15 ProductAccessTokenRequestItem, SuotarCallContext, SuotarEndpoint, SuotarItemStatus,
16};
17
18use super::{CreditRegistrationPhase, PhaseContext, PhaseScope, every_item_failed_transiently};
19
20pub async fn run(ctx: &PhaseContext<'_>, scope: &PhaseScope) -> anyhow::Result<PhaseRunOutcome> {
21 let endpoint = SuotarEndpoint::ProductAccessTokens;
22 let mut conn = ctx.pool.acquire().await?;
23 let products = get_stalest_product_ids_for_enabled_modules(
24 &mut conn,
25 endpoint.max_batch_size() as i64,
26 scope.course_id,
27 )
28 .await?;
29 let attempted = i32::try_from(products.len()).unwrap_or(i32::MAX);
30 if products.is_empty() {
31 return Ok(PhaseRunOutcome::default());
32 }
33
34 let items: Vec<_> = products
35 .iter()
36 .map(|product_id| ProductAccessTokenRequestItem {
37 request_item_id: request_item_id(product_id),
38 open_university_product_id: product_id.clone(),
39 })
40 .collect();
41 drop(conn);
43 let response = ctx
44 .suotar_client
45 .resolve_product_access_tokens(
46 SuotarCallContext::new(ctx.worker_name(CreditRegistrationPhase::ProductTokenRefresh)),
47 items,
48 )
49 .await;
50 let response = match response {
51 Ok(response) => response,
52 Err(error) => {
53 let mut conn = ctx.pool.acquire().await?;
54 let message = scrub_text(error.message());
55 for product_id in &products {
56 record_refresh_failure(&mut conn, product_id, &message).await?;
57 }
58 return Ok(PhaseRunOutcome {
59 items_processed: attempted,
60 items_failed: attempted,
61 error: Some(message),
62 });
63 }
64 };
65
66 let mut conn = ctx.pool.acquire().await?;
67 let mut items_failed = 0;
68 for product_id in &products {
69 let item = response.item(&request_item_id(product_id));
70 let refreshed = match item {
71 Some(item) if item.status == SuotarItemStatus::Ok => item.result.as_ref(),
72 Some(item) => {
73 let message = item
74 .error
75 .as_ref()
76 .map(|error| scrub_text(&error.message))
77 .unwrap_or_else(|| item.code.clone());
78 record_refresh_failure(&mut conn, product_id, &message).await?;
79 items_failed += 1;
80 continue;
81 }
82 None => {
83 record_refresh_failure(
84 &mut conn,
85 product_id,
86 "The study registry did not answer for this product.",
87 )
88 .await?;
89 items_failed += 1;
90 continue;
91 }
92 };
93 let Some(refreshed) = refreshed else {
94 record_refresh_failure(
95 &mut conn,
96 product_id,
97 "The study registry reported success but sent no token.",
98 )
99 .await?;
100 items_failed += 1;
101 continue;
102 };
103 upsert(
104 &mut conn,
105 &NewOpenUniversityProductAccessToken {
106 open_university_product_id: product_id.clone(),
107 access_token: DbSecret::new(refreshed.access_token.clone()),
108 state: refreshed.state.clone(),
109 document_state: refreshed.document_state.clone(),
110 suotar_token_id: Some(refreshed.id.clone()),
111 },
112 )
113 .await?;
114 }
115
116 Ok(PhaseRunOutcome {
117 items_processed: attempted,
118 items_failed,
119 error: every_item_failed_transiently(&response)
120 .then(|| "Every product of the batch came back transiently unavailable.".to_string()),
121 })
122}
123
124fn request_item_id(open_university_product_id: &str) -> String {
127 format!("oup-{open_university_product_id}")
128}