Skip to main content

headless_lms_credit_registration/runtime/suotar/
executor.rs

1//! Sending one batch of a worker flow: fresh request ids, the limiter, the call, the gate record,
2//! and pairing each row with its answer and audit.
3
4use 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
19/// Sends `entries` in one request to `E`, every item under a fresh requestItemId, a resent half's
20/// included.
21pub(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            // Suotar validates every item before acting on any, so one bad row takes its whole
63            // batch down with a malformed-request refusal, and that refusal proves nothing was
64            // acted on.
65            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
143/// What the audit needs of one item after the items themselves went into the request.
144struct SentItem {
145    request_item_id: String,
146    student_number: Option<StudentNumber>,
147}
148
149/// The request bodies as sent, kept alongside the typed items so a rejected batch can pair each row
150/// with what was actually asked of it for the audit log.
151fn 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
158/// The response item for one request item, read from the raw body rather than rebuilt from the
159/// typed value, so the audit trail holds what actually arrived.
160fn 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}