headless_lms_credit_registration/runtime/suotar/
executor.rs1use chrono::Utc;
5use headless_lms_models::credit_registration_events::SuotarAnswer;
6use headless_lms_utils::services::suotar::{
7 BatchEndpoint, SuotarErrorVariant, SuotarRequestItem, SuotarResponseItem, new_request_item_id,
8};
9use tracing::Instrument;
10
11use super::SuotarStudyRegistry;
12use super::decode::registry_error;
13use super::gate::Exchange;
14use crate::registry::{
15 AnsweredRow, BatchEntry, BatchOptions, BatchReply, ExchangeAudit, RefusedFor, RefusedRow,
16 StudentNumber,
17};
18
19pub(super) async fn send_batch<E: BatchEndpoint, K, R, A>(
22 registry: &mut SuotarStudyRegistry<'_>,
23 entries: Vec<BatchEntry<K, R>>,
24 options: BatchOptions,
25 encode: impl Fn(&R, String) -> E::Item,
26 decode: impl Fn(&SuotarResponseItem<E::Result>) -> A,
27) -> BatchReply<K, R, A> {
28 let endpoint = E::ENDPOINT;
29 let items: Vec<E::Item> = entries
30 .iter()
31 .map(|entry| encode(&entry.request, new_request_item_id()))
32 .collect();
33 let requests = requests_json(&items);
34 let sent_items: Vec<SentItem> = items
35 .iter()
36 .map(|item| SentItem {
37 request_item_id: item.request_item_id().to_string(),
38 student_number: item.student_number().cloned().map(StudentNumber::new),
39 })
40 .collect();
41 let item_count = items.len();
42 if options.is_resent_half {
43 registry.gate.spend_split(endpoint, item_count);
44 } else {
45 registry.gate.spend(endpoint, item_count);
46 }
47 let context = registry.call_context(options.registration_ids);
48 let span = super::request_span(endpoint, item_count);
49 if options.is_resent_half {
50 span.record("resent_half", true);
51 }
52 let requested_at = Utc::now();
53 let sent = registry
54 .client
55 .post::<E>(context, items)
56 .instrument(span)
57 .await;
58 let answered_at = Utc::now();
59 let response = match sent {
60 Ok(response) => response,
61 Err(error) => {
62 let is_isolated =
66 error.variant == SuotarErrorVariant::MalformedRequest && options.may_split;
67 if is_isolated && item_count > 1 {
68 return BatchReply::RefusedAsMalformed {
69 entries,
70 error: registry_error(&error),
71 };
72 }
73 let refused_for = if is_isolated {
74 registry
75 .gate
76 .record(endpoint, Exchange::RefusedAlone(&error));
77 RefusedFor::RowAlone
78 } else {
79 registry.gate.record(endpoint, Exchange::Refused(&error));
80 RefusedFor::WholeBatch
81 };
82 let rows = entries
83 .into_iter()
84 .zip(sent_items)
85 .zip(requests)
86 .map(|((entry, sent), request)| RefusedRow {
87 row: entry.row,
88 audit: ExchangeAudit {
89 call_id: None,
90 endpoint,
91 requested_at,
92 answered_at,
93 answer: SuotarAnswer::Refused,
94 request_item_id: sent.request_item_id,
95 request,
96 response: None,
97 sent_student_number: None,
98 },
99 })
100 .collect();
101 return BatchReply::Refused {
102 rows,
103 error: registry_error(&error),
104 refused_for,
105 };
106 }
107 };
108 registry.gate.record(
109 endpoint,
110 Exchange::answered(&response, options.all_unavailable_error),
111 );
112 let rows = entries
113 .into_iter()
114 .zip(sent_items)
115 .zip(requests)
116 .map(|((entry, sent), request)| {
117 let request_item_id = sent.request_item_id;
118 let item = response.item(&request_item_id);
119 AnsweredRow {
120 row: entry.row,
121 answer: item.map(&decode),
122 audit: ExchangeAudit {
123 call_id: response.call_id,
124 endpoint,
125 requested_at,
126 answered_at,
127 answer: if item.is_some() {
128 SuotarAnswer::Answered
129 } else {
130 SuotarAnswer::Unanswered
131 },
132 response: response_item_json(&response.raw_response, &request_item_id),
133 request_item_id,
134 request,
135 sent_student_number: sent.student_number,
136 },
137 }
138 })
139 .collect();
140 BatchReply::Answered(rows)
141}
142
143struct SentItem {
145 request_item_id: String,
146 student_number: Option<StudentNumber>,
147}
148
149fn requests_json<T: serde::Serialize>(items: &[T]) -> Vec<serde_json::Value> {
152 items
153 .iter()
154 .map(|item| serde_json::to_value(item).unwrap_or_default())
155 .collect()
156}
157
158fn response_item_json(
161 raw_response: &serde_json::Value,
162 request_item_id: &str,
163) -> Option<serde_json::Value> {
164 raw_response
165 .as_array()?
166 .iter()
167 .find(|item| item.get("requestItemId").and_then(|id| id.as_str()) == Some(request_item_id))
168 .cloned()
169}
170
171#[cfg(test)]
172mod tests {
173 use super::*;
174
175 #[test]
176 fn a_response_item_is_found_by_its_request_item_id() {
177 let raw = serde_json::json!([
178 { "requestItemId": "item-1", "status": "ok", "code": "sent" },
179 { "requestItemId": "item-2", "status": "error", "code": "sisuTimeout" },
180 ]);
181 assert_eq!(
182 response_item_json(&raw, "item-2").and_then(|item| item
183 .get("code")
184 .and_then(|code| code.as_str().map(str::to_string))),
185 Some("sisuTimeout".to_string())
186 );
187 assert_eq!(response_item_json(&raw, "item-9"), None);
188 }
189}