headless_lms_credit_registration/runtime/suotar/
mod.rs1mod breaker;
5mod codes;
6#[cfg(test)]
7mod contract_tests;
8mod course_codes;
9mod decode;
10mod encode;
11mod executor;
12mod gate;
13mod health_report;
14mod rate_limit;
15mod rosters;
16
17pub use breaker::is_waiting_to_probe;
18pub use codes::is_waiting_item;
19pub(super) use health_report::{report_breakers, report_rate_limits};
20
21use headless_lms_models::suotar_api_calls::SuotarEndpoint;
22use headless_lms_models::suotar_circuit_breakers::BreakerTarget;
23use headless_lms_utils::services::suotar::{
24 INTERACTIVE_REQUEST_TIMEOUT, SuotarCallContext, SuotarClient, endpoints, new_request_item_id,
25};
26use itertools::Itertools;
27use tracing::{Instrument, Span};
28use uuid::Uuid;
29
30use crate::phase::CreditRegistrationPhase;
31use crate::registry::{
32 AttainmentSubmission, BatchEntry, BatchOptions, BatchReply, CourseCode, CourseCodeVerdicts,
33 EnrolmentAnswer, EnrolmentLookup, ImportAnswer, InteractiveStudyRegistry, PersonAnswer,
34 PersonLookup, PersonLookupError, RegistryError, RegistryOperation, RegistryPerson, RosterCode,
35 RosterListing, RosterSearch, StudentNumber, StudyRegistry, VerificationAnswer,
36 VerificationRequest,
37};
38use crate::runtime::PhaseSkipReason;
39use crate::runtime::process_local::ScopeKey;
40use gate::StudyRegistryGate;
41use headless_lms_models::credit_registrations::RegistrationScope;
42
43pub fn reset_rate_limits(scope: &RegistrationScope) {
45 rate_limit::reset(&ScopeKey::of(scope));
46}
47
48pub(super) fn endpoint_of(operation: RegistryOperation) -> SuotarEndpoint {
51 match operation {
52 RegistryOperation::ResolvePersons => SuotarEndpoint::ResolvePersons,
53 RegistryOperation::ResolveEnrolments => SuotarEndpoint::ResolveEnrolments,
54 RegistryOperation::ImportAttainments => SuotarEndpoint::ImportAttainments,
55 RegistryOperation::VerifyAttainments => SuotarEndpoint::VerifyAttainments,
56 RegistryOperation::ListCourseRoster => SuotarEndpoint::ListByCourse,
57 RegistryOperation::ValidateCourseCodes => SuotarEndpoint::ValidateCourseCodes,
58 }
59}
60
61fn study_registry_endpoints(
64 phase: CreditRegistrationPhase,
65) -> impl Iterator<Item = SuotarEndpoint> {
66 phase
67 .spec()
68 .operations
69 .iter()
70 .map(|&operation| endpoint_of(operation))
71}
72
73pub fn endpoints_paused_by(process_name: &str, target: BreakerTarget) -> Vec<SuotarEndpoint> {
75 CreditRegistrationPhase::ALL
76 .into_iter()
77 .filter(|phase| {
78 let spec = phase.spec();
79 spec.process.as_str() == process_name && spec.breakers.contains(&target)
80 })
81 .flat_map(study_registry_endpoints)
82 .unique()
83 .collect()
84}
85
86pub fn max_study_registry_wait(phase: CreditRegistrationPhase) -> std::time::Duration {
88 study_registry_endpoints(phase)
89 .map(SuotarEndpoint::request_timeout)
90 .sum()
91}
92
93fn request_span(endpoint: SuotarEndpoint, items: usize) -> Span {
96 debug_span!(
97 "study_registry_request",
98 ?endpoint,
99 items,
100 resent_half = false
101 )
102}
103
104pub(super) struct SuotarStudyRegistry<'a> {
107 client: &'a SuotarClient,
108 worker_name: String,
110 gate: StudyRegistryGate,
111 #[cfg(test)]
113 request_timeout: Option<std::time::Duration>,
114}
115
116impl<'a> SuotarStudyRegistry<'a> {
117 pub(super) fn admit(
119 client: &'a SuotarClient,
120 worker_name: String,
121 phase: CreditRegistrationPhase,
122 scope: &RegistrationScope,
123 test_mode: bool,
124 ) -> Result<Self, PhaseSkipReason> {
125 Ok(Self {
126 client,
127 worker_name,
128 gate: StudyRegistryGate::admit(phase, scope, test_mode)?,
129 #[cfg(test)]
130 request_timeout: None,
131 })
132 }
133
134 pub(super) fn finish(self) -> Option<String> {
137 self.gate.settle()
138 }
139
140 fn call_context(&self, registration_ids: Vec<Uuid>) -> SuotarCallContext {
141 SuotarCallContext {
142 worker_name: self.worker_name.clone(),
143 credit_registration_ids: registration_ids,
144 #[cfg(test)]
145 request_timeout: self.request_timeout,
146 #[cfg(not(test))]
147 request_timeout: None,
148 }
149 }
150}
151
152impl StudyRegistry for SuotarStudyRegistry<'_> {
153 fn allowance(&self, operation: RegistryOperation) -> usize {
154 self.gate.allowance(endpoint_of(operation))
155 }
156
157 fn roster_request_size(&self) -> usize {
158 if self.gate.is_probe() {
159 1
160 } else {
161 SuotarEndpoint::ListByCourse.max_batch_size()
162 }
163 }
164
165 async fn resolve_persons<K>(
166 &mut self,
167 entries: Vec<BatchEntry<K, PersonLookup>>,
168 options: BatchOptions,
169 ) -> BatchReply<K, PersonLookup, PersonAnswer> {
170 executor::send_batch::<endpoints::ResolvePersons, _, _, _>(
171 self,
172 entries,
173 options,
174 encode::person_lookup_item,
175 decode::person_answer,
176 )
177 .await
178 }
179
180 async fn resolve_enrolments<K>(
181 &mut self,
182 entries: Vec<BatchEntry<K, EnrolmentLookup>>,
183 options: BatchOptions,
184 ) -> BatchReply<K, EnrolmentLookup, EnrolmentAnswer> {
185 executor::send_batch::<endpoints::ResolveEnrolments, _, _, _>(
186 self,
187 entries,
188 options,
189 encode::enrolment_item,
190 decode::enrolment_answer,
191 )
192 .await
193 }
194
195 async fn import_attainments<K>(
196 &mut self,
197 entries: Vec<BatchEntry<K, AttainmentSubmission>>,
198 options: BatchOptions,
199 ) -> BatchReply<K, AttainmentSubmission, ImportAnswer> {
200 executor::send_batch::<endpoints::ImportAttainments, _, _, _>(
201 self,
202 entries,
203 options,
204 encode::import_item,
205 decode::import_answer,
206 )
207 .await
208 }
209
210 async fn verify_attainments<K>(
211 &mut self,
212 entries: Vec<BatchEntry<K, VerificationRequest>>,
213 options: BatchOptions,
214 ) -> BatchReply<K, VerificationRequest, VerificationAnswer> {
215 executor::send_batch::<endpoints::VerifyAttainments, _, _, _>(
216 self,
217 entries,
218 options,
219 encode::verify_item,
220 decode::verification_answer,
221 )
222 .await
223 }
224
225 async fn list_course_roster(
226 &mut self,
227 request: &[RosterCode],
228 ) -> Result<RosterListing, RegistryError> {
229 rosters::list(self, request).await
230 }
231
232 async fn validate_course_codes(
233 &mut self,
234 codes: &[CourseCode],
235 ) -> Result<CourseCodeVerdicts, RegistryError> {
236 course_codes::validate(self, codes).await
237 }
238}
239
240pub(super) struct InteractiveSuotar<'a> {
243 client: &'a SuotarClient,
244 worker_name: String,
246}
247
248impl<'a> InteractiveSuotar<'a> {
249 pub(super) fn new(client: &'a SuotarClient, worker_name: String) -> Self {
250 Self {
251 client,
252 worker_name,
253 }
254 }
255
256 fn call_context(&self) -> SuotarCallContext {
257 SuotarCallContext {
258 worker_name: self.worker_name.clone(),
259 credit_registration_ids: Vec::new(),
260 request_timeout: Some(INTERACTIVE_REQUEST_TIMEOUT),
261 }
262 }
263}
264
265impl InteractiveStudyRegistry for InteractiveSuotar<'_> {
266 async fn look_up_person(
267 &self,
268 student_number: &StudentNumber,
269 ) -> Result<Option<RegistryPerson>, PersonLookupError> {
270 let request_item_id = new_request_item_id();
271 let item = encode::person_item(student_number, request_item_id.clone());
272 let span = request_span(SuotarEndpoint::ResolvePersons, 1);
273 let response = self
274 .client
275 .post::<endpoints::ResolvePersons>(self.call_context(), vec![item])
276 .instrument(span)
277 .await
278 .map_err(|error| {
279 warn!(error = %error, "Could not resolve a student number in the study registry");
280 PersonLookupError::StudyRegistryUnavailable
281 })?;
282 match response.item(&request_item_id) {
283 Some(item) => decode::person_lookup(item),
284 None => Err(PersonLookupError::ItemMissingFromResponse),
285 }
286 }
287
288 async fn search_course_rosters(
289 &self,
290 codes: &[CourseCode],
291 student_number: &StudentNumber,
292 ) -> RosterSearch {
293 rosters::search(self, codes, student_number).await
294 }
295}