Skip to main content

headless_lms_credit_registration/runtime/suotar/
gate.rs

1//! The one place an iteration's Suotar requests meet the circuit breakers and the limiter: every
2//! request is spent through the gate and recorded in it, and [`StudyRegistryGate::settle`] turns
3//! the records into the breaker and limiter effects.
4
5use headless_lms_models::library::credit_registration::scrub::scrub_text;
6use headless_lms_utils::services::suotar::{SuotarBatchResponse, SuotarEndpoint, SuotarError};
7
8use super::breaker::{self, BreakerTarget};
9use super::codes::{is_all_unavailable, is_only_sisu_timeouts};
10use super::decode::registry_error;
11use super::rate_limit;
12use crate::phase::CreditRegistrationPhase;
13use crate::runtime::PhaseSkipReason;
14use crate::runtime::process_local::ScopeKey;
15use headless_lms_models::credit_registrations::RegistrationScope;
16
17/// What one request to Suotar came to.
18pub(super) enum Exchange<'a> {
19    /// Suotar answered; `unavailable` when every item came back unavailable.
20    Answered { unavailable: Option<Unavailable> },
21    /// Suotar refused the whole request, or it never got there.
22    Refused(&'a SuotarError),
23    /// Suotar refused a request that only one row or code of ours can be blamed for — says
24    /// nothing about Suotar's own health.
25    RefusedAlone(&'a SuotarError),
26}
27
28/// A batch whose every item came back unavailable, which fails the iteration.
29pub(super) struct Unavailable {
30    /// The iteration's error for it.
31    pub message: &'static str,
32    /// Sisu timed out on every submission: Suotar itself answered, so only the phase that submits
33    /// pauses.
34    pub is_only_sisu_timeouts: bool,
35}
36
37impl Exchange<'_> {
38    /// An answer; `unavailable_message` is the iteration's error if every item came back
39    /// unavailable.
40    pub(super) fn answered<R>(
41        response: &SuotarBatchResponse<R>,
42        unavailable_message: &'static str,
43    ) -> Self {
44        let endpoint = response.endpoint;
45        let items = &response.items;
46        let answers = items.iter().map(|item| (item.status, item.code.as_str()));
47        Self::Answered {
48            unavailable: is_all_unavailable(endpoint, answers).then(|| Unavailable {
49                message: unavailable_message,
50                is_only_sisu_timeouts: is_only_sisu_timeouts(
51                    endpoint,
52                    items.iter().map(|item| item.code.as_str()),
53                ),
54            }),
55        }
56    }
57}
58
59/// What the exchanges of one iteration said about the study registry.
60#[derive(Debug, Default)]
61struct Tally {
62    has_answer: bool,
63    /// Suotar failing in a way that can pass: a transport error, a timeout or a 5xx, or every item
64    /// of an answer unavailable.
65    has_registry_failure: bool,
66    has_sisu_outage: bool,
67    /// The first failure of the iteration, which stands for it.
68    error: Option<String>,
69    /// The first request refused with one row or code alone in it.
70    isolated: Option<String>,
71    /// The endpoints whose own failures drop their limiter to its floor.
72    failed_endpoints: Vec<SuotarEndpoint>,
73}
74
75/// What [`StudyRegistryGate::settle`] makes of a [`Tally`].
76#[derive(Debug, Clone, Copy, PartialEq, Eq)]
77enum Verdict {
78    /// Nothing reached Suotar, or nothing that says whether it's up: counts against no breaker and
79    /// clears no failure run — an empty queue is common, and clearing on it would hide a real outage.
80    Idle,
81    Healthy,
82    /// Suotar answered, and Sisu timed out on every submission.
83    SisuOutage,
84    RegistryDown,
85}
86
87fn verdict(tally: &Tally) -> Verdict {
88    if tally.has_registry_failure {
89        Verdict::RegistryDown
90    } else if tally.has_sisu_outage {
91        Verdict::SisuOutage
92    } else if tally.has_answer {
93        Verdict::Healthy
94    } else {
95        Verdict::Idle
96    }
97}
98
99/// One iteration's passage through the circuit breakers and the limiter. The adapter asks it how
100/// much it may send, spends that before sending, and records what came back; nothing else in an
101/// iteration touches [`breaker`] or [`rate_limit`].
102pub(super) struct StudyRegistryGate {
103    key: ScopeKey,
104    phase: CreditRegistrationPhase,
105    test_mode: bool,
106    /// After a breaker's cooldown: the iteration may send one single-item request, whichever of its
107    /// flows sends it, and the breaker closes only if that request succeeds.
108    is_probe: bool,
109    has_probed: bool,
110    tally: Tally,
111}
112
113impl StudyRegistryGate {
114    /// Lets the iteration run, or says why it waits. Only a phase that calls the study registry
115    /// ever waits: an outage must not stall the database-only phases.
116    pub(super) fn admit(
117        phase: CreditRegistrationPhase,
118        scope: &RegistrationScope,
119        test_mode: bool,
120    ) -> Result<Self, PhaseSkipReason> {
121        let key = ScopeKey::of(scope);
122        let targets = phase.spec().breakers;
123        if targets.iter().any(|&target| breaker::is_open(&key, target)) {
124            return Err(PhaseSkipReason::CircuitBreakerOpen);
125        }
126        let is_probe = targets
127            .iter()
128            .any(|&target| breaker::is_half_open(&key, target));
129        Ok(Self {
130            key,
131            phase,
132            test_mode,
133            is_probe,
134            has_probed: false,
135            tally: Tally::default(),
136        })
137    }
138
139    /// Whether the iteration is a breaker's probe: see [`Self::allowance`].
140    pub(super) fn is_probe(&self) -> bool {
141        self.is_probe
142    }
143
144    /// How many items one request to `endpoint` may carry now: its batch size, cut to what the
145    /// limiter allows, or during a probe one item for the first request and none after it. For
146    /// `list-by-course`, whose limiter counts requests, how many requests.
147    pub(super) fn allowance(&self, endpoint: SuotarEndpoint) -> usize {
148        if self.is_probe {
149            let allowance = usize::from(!self.has_probed);
150            if breaker::is_half_open(&self.key, BreakerTarget::StudyRegistry) {
151                debug!(
152                    ?endpoint,
153                    allowance, "Study registry breaker is half-open; probing with one item"
154                );
155            } else {
156                debug!(
157                    ?endpoint,
158                    allowance, "Sisu submissions breaker is half-open; probing with one item"
159                );
160            }
161            return allowance;
162        }
163        let limit = endpoint
164            .max_batch_size()
165            .min(rate_limit::available(&self.key, endpoint));
166        trace!(
167            ?endpoint,
168            limit, "Computed the claim limit for a Suotar endpoint"
169        );
170        limit
171    }
172
173    /// Spends `count` of what [`Self::allowance`] allowed, just before the request leaves.
174    pub(super) fn spend(&mut self, endpoint: SuotarEndpoint, count: usize) {
175        rate_limit::take(&self.key, endpoint, count);
176        self.has_probed = true;
177    }
178
179    /// [`Self::spend`] for the resent halves of a batch refused as malformed, which may go past
180    /// the allowance; see [`rate_limit::overdraw`].
181    pub(super) fn spend_split(&mut self, endpoint: SuotarEndpoint, count: usize) {
182        rate_limit::overdraw(&self.key, endpoint, count);
183        self.has_probed = true;
184    }
185
186    /// Records what one request to `endpoint` came to, for [`Self::settle`].
187    pub(super) fn record(&mut self, endpoint: SuotarEndpoint, exchange: Exchange<'_>) {
188        let tally = &mut self.tally;
189        match exchange {
190            Exchange::Answered { unavailable: None } => tally.has_answer = true,
191            Exchange::Answered {
192                unavailable: Some(unavailable),
193            } => {
194                tally.has_answer = true;
195                tally
196                    .error
197                    .get_or_insert_with(|| unavailable.message.to_string());
198                if unavailable.is_only_sisu_timeouts {
199                    tally.has_sisu_outage = true;
200                } else {
201                    tally.has_registry_failure = true;
202                    tally.failed_endpoints.push(endpoint);
203                }
204            }
205            Exchange::Refused(error) => {
206                tally
207                    .error
208                    .get_or_insert_with(|| scrub_text(error.message()));
209                // Only Suotar failing counts: not our own request or credentials, and not a request
210                // that never left.
211                tally.has_registry_failure |=
212                    error.was_sent && registry_error(error).kind.is_outage();
213                tally.failed_endpoints.push(endpoint);
214            }
215            Exchange::RefusedAlone(error) => {
216                tally
217                    .isolated
218                    .get_or_insert_with(|| scrub_text(error.message()));
219            }
220        }
221    }
222
223    /// Applies what the iteration's exchanges said to the breakers and the limiter, and returns the
224    /// iteration's error, if any. A request refused with its row alone is reported only when
225    /// nothing answered, and counts against no breaker.
226    ///
227    /// The limiter drops to its floor whenever the shared breaker trips or closes again, so the ramp
228    /// back starts from the probe that got through rather than from a failure a long cooldown ago.
229    pub(super) fn settle(self) -> Option<String> {
230        let key = &self.key;
231        let phase = self.phase.as_str();
232        let base_cooldown = breaker::cooldown(self.test_mode);
233        let submits_to_sisu = self
234            .phase
235            .spec()
236            .breakers
237            .contains(&BreakerTarget::SisuSubmissions);
238        rate_limit::drop_to_floor(key, &self.tally.failed_endpoints);
239        match verdict(&self.tally) {
240            Verdict::Idle => {}
241            Verdict::Healthy => {
242                record_study_registry_success(key);
243                if submits_to_sisu
244                    && let Some(trip_count) =
245                        breaker::record_success(key, BreakerTarget::SisuSubmissions)
246                {
247                    info!(phase, trip_count, "Sisu submissions circuit breaker closed");
248                }
249            }
250            Verdict::SisuOutage => {
251                record_study_registry_success(key);
252                if let Some(trip) =
253                    breaker::record_failure(key, BreakerTarget::SisuSubmissions, base_cooldown)
254                {
255                    warn!(
256                        phase,
257                        cooldown_secs = trip.cooldown.as_secs(),
258                        consecutive_failures = trip.consecutive_failures,
259                        trip_count = trip.trip_count,
260                        "Pausing phase after consecutive Sisu timeouts"
261                    );
262                }
263            }
264            Verdict::RegistryDown => {
265                if let Some(trip) =
266                    breaker::record_failure(key, BreakerTarget::StudyRegistry, base_cooldown)
267                {
268                    rate_limit::drop_to_floor(key, &rate_limit::LIMITED_ENDPOINTS);
269                    warn!(
270                        phase,
271                        cooldown_secs = trip.cooldown.as_secs(),
272                        consecutive_failures = trip.consecutive_failures,
273                        trip_count = trip.trip_count,
274                        "Pausing study registry phases after consecutive failures"
275                    );
276                }
277            }
278        }
279        let tally = self.tally;
280        tally.error.or(if tally.has_answer {
281            None
282        } else {
283            tally.isolated
284        })
285    }
286}
287
288fn record_study_registry_success(key: &ScopeKey) {
289    if let Some(trip_count) = breaker::record_success(key, BreakerTarget::StudyRegistry) {
290        rate_limit::drop_to_floor(key, &rate_limit::LIMITED_ENDPOINTS);
291        info!(trip_count, "Study registry circuit breaker closed");
292    }
293}