Skip to main content

headless_lms_credit_registration/runtime/suotar/
mod.rs

1//! The study registry over Suotar, for phase iterations and for manual actions: the only part of
2//! the pipeline that names Suotar's wire items and codes, or sends through its client.
3
4mod 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
43/// Forgets the limiter state of `scope`, so its endpoints are back at full rate with a full burst.
44pub fn reset_rate_limits(scope: &RegistrationScope) {
45    rate_limit::reset(&ScopeKey::of(scope));
46}
47
48/// The Suotar endpoint that serves `operation`, which is also what the call log, the limiter and the
49/// dashboard key on.
50pub(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
61/// The Suotar endpoints one iteration of `phase` calls, in order: its
62/// [`crate::PhaseSpec::operations`] as the call log and the dashboard name them.
63fn 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
73/// The Suotar endpoints whose phases `process_name`'s `target` breaker pauses, each once.
74pub 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
86/// The longest one iteration of `phase` may wait on the study registry before its calls time out.
87pub fn max_study_registry_wait(phase: CreditRegistrationPhase) -> std::time::Duration {
88    study_registry_endpoints(phase)
89        .map(SuotarEndpoint::request_timeout)
90        .sum()
91}
92
93/// The span every request to `endpoint` runs in, whichever of the adapter's call sites sends it.
94/// A resent half records `resent_half` on it.
95fn 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
104/// The [`StudyRegistry`] over Suotar for one phase iteration: every request is spent through, and
105/// recorded in, the iteration's gate.
106pub(super) struct SuotarStudyRegistry<'a> {
107    client: &'a SuotarClient,
108    /// The audit log's `worker_name` for every call.
109    worker_name: String,
110    gate: StudyRegistryGate,
111    /// Replaces the endpoints' own timeouts, which run to minutes.
112    #[cfg(test)]
113    request_timeout: Option<std::time::Duration>,
114}
115
116impl<'a> SuotarStudyRegistry<'a> {
117    /// The registry one phase iteration sends through, or why the iteration waits.
118    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    /// Applies what the iteration's calls said to the breakers and the limiter, and returns the
135    /// iteration's error, if any.
136    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
240/// The [`InteractiveStudyRegistry`] over Suotar, for one manual action: the interactive timeout,
241/// and no gate.
242pub(super) struct InteractiveSuotar<'a> {
243    client: &'a SuotarClient,
244    /// The audit log's `worker_name` for every call.
245    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}