headless_lms_credit_registration/runtime/suotar/
gate.rs1use 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
17pub(super) enum Exchange<'a> {
19 Answered { unavailable: Option<Unavailable> },
21 Refused(&'a SuotarError),
23 RefusedAlone(&'a SuotarError),
26}
27
28pub(super) struct Unavailable {
30 pub message: &'static str,
32 pub is_only_sisu_timeouts: bool,
35}
36
37impl Exchange<'_> {
38 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#[derive(Debug, Default)]
61struct Tally {
62 has_answer: bool,
63 has_registry_failure: bool,
66 has_sisu_outage: bool,
67 error: Option<String>,
69 isolated: Option<String>,
71 failed_endpoints: Vec<SuotarEndpoint>,
73}
74
75#[derive(Debug, Clone, Copy, PartialEq, Eq)]
77enum Verdict {
78 Idle,
81 Healthy,
82 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
99pub(super) struct StudyRegistryGate {
103 key: ScopeKey,
104 phase: CreditRegistrationPhase,
105 test_mode: bool,
106 is_probe: bool,
109 has_probed: bool,
110 tally: Tally,
111}
112
113impl StudyRegistryGate {
114 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 pub(super) fn is_probe(&self) -> bool {
141 self.is_probe
142 }
143
144 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 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 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 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 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 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}