headless_lms_credit_registration/runtime/suotar/
breaker.rs1use std::time::{Duration, Instant};
14
15use chrono::{DateTime, TimeDelta, Utc};
16
17use headless_lms_models::suotar_api_calls::SuotarEndpoint;
18
19use crate::runtime::process_local::{LastReported, ProcessLocalMap, ScopeKey};
20
21const MAX_CONSECUTIVE_SUOTAR_FAILURES: u32 = 5;
22const SUOTAR_COOLDOWN: Duration = Duration::from_secs(5 * 60);
25const MAX_COOLDOWN_TRIPS: u32 = 3;
26const TEST_SUOTAR_COOLDOWN: Duration = Duration::from_secs(5);
29
30pub(super) use headless_lms_models::suotar_circuit_breakers::BreakerTarget;
31
32const FAILURE_RUN_MEMORY: Duration = Duration::from_secs(2 * 60 * 60);
36const _: () = assert!(
37 FAILURE_RUN_MEMORY.as_secs()
38 >= 2 * SuotarEndpoint::ImportAttainments
39 .request_timeout()
40 .as_secs()
41);
42
43#[derive(Debug, Clone)]
44struct BreakerState {
45 consecutive_failures: u32,
46 open_until: Option<Instant>,
47 last_failure_at: Instant,
48 trip_count: u32,
50}
51
52impl BreakerState {
53 fn is_live(&self, now: Instant) -> bool {
55 self.open_until.is_some_and(|until| now < until)
56 || now.duration_since(self.last_failure_at) < FAILURE_RUN_MEMORY
57 }
58}
59
60type BreakerKey = (ScopeKey, BreakerTarget);
61
62static BREAKERS: ProcessLocalMap<BreakerKey, BreakerState> = ProcessLocalMap::new();
63
64pub(super) static REPORTED: LastReported<BreakerTarget, BreakerSnapshot> = LastReported::new();
66
67pub(super) fn cooldown(test_mode: bool) -> Duration {
69 if test_mode {
70 TEST_SUOTAR_COOLDOWN
71 } else {
72 SUOTAR_COOLDOWN
73 }
74}
75
76pub(super) fn is_open(scope: &ScopeKey, target: BreakerTarget) -> bool {
78 let key = (scope.clone(), target);
79 let now = Instant::now();
80 let mut breakers = BREAKERS.lock();
81 let Some(state) = breakers.get(&key) else {
82 return false;
83 };
84 if state.open_until.is_some_and(|until| now < until) {
85 return true;
86 }
87 if state.is_live(now) {
88 return false;
89 }
90 breakers.remove(&key);
93 false
94}
95
96pub(super) fn is_half_open(scope: &ScopeKey, target: BreakerTarget) -> bool {
99 let now = Instant::now();
100 BREAKERS
101 .lock()
102 .get(&(scope.clone(), target))
103 .is_some_and(|state| {
104 state.is_live(now)
105 && is_waiting_to_probe(state.consecutive_failures, state.open_until, now)
106 })
107}
108
109pub fn is_waiting_to_probe<T: PartialOrd>(
113 consecutive_failures: u32,
114 open_until: Option<T>,
115 now: T,
116) -> bool {
117 open_until.is_some_and(|until| now >= until)
118 && consecutive_failures >= MAX_CONSECUTIVE_SUOTAR_FAILURES
119}
120
121#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
123pub(super) struct BreakerSnapshot {
124 pub open: bool,
125 pub consecutive_failures: u32,
126 pub open_until: Option<DateTime<Utc>>,
128 pub trip_count: u32,
129}
130
131impl BreakerSnapshot {
132 pub(super) fn is_same_report(&self, other: &Self) -> bool {
135 let to_second = |snapshot: &Self| snapshot.open_until.map(|until| until.timestamp());
136 self.consecutive_failures == other.consecutive_failures
137 && self.trip_count == other.trip_count
138 && to_second(self) == to_second(other)
139 }
140}
141
142pub(super) fn snapshot(scope: &ScopeKey, target: BreakerTarget) -> BreakerSnapshot {
145 let breakers = BREAKERS.lock();
146 let now = Instant::now();
147 let Some(state) = breakers
148 .get(&(scope.clone(), target))
149 .filter(|state| state.is_live(now))
150 else {
151 return BreakerSnapshot::default();
152 };
153 let wall_now = Utc::now();
154 BreakerSnapshot {
155 open: state.open_until.is_some_and(|until| now < until),
156 consecutive_failures: state.consecutive_failures,
157 open_until: state
158 .open_until
159 .map(|until| match until.checked_duration_since(now) {
160 Some(left) => wall_now + TimeDelta::from_std(left).unwrap_or_default(),
161 None => {
162 wall_now - TimeDelta::from_std(now.duration_since(until)).unwrap_or_default()
163 }
164 }),
165 trip_count: state.trip_count,
166 }
167}
168
169pub(super) fn record_success(scope: &ScopeKey, target: BreakerTarget) -> Option<u32> {
171 BREAKERS
172 .lock()
173 .remove(&(scope.clone(), target))
174 .map(|state| state.trip_count)
175 .filter(|&trip_count| trip_count > 0)
176}
177
178pub(super) struct BreakerTrip {
180 pub cooldown: Duration,
181 pub consecutive_failures: u32,
182 pub trip_count: u32,
183}
184
185pub(super) fn record_failure(
188 scope: &ScopeKey,
189 target: BreakerTarget,
190 base_cooldown: Duration,
191) -> Option<BreakerTrip> {
192 let now = Instant::now();
193 let mut breakers = BREAKERS.lock();
194 breakers.retain(|_, state| state.is_live(now));
195 let state = breakers
196 .entry((scope.clone(), target))
197 .or_insert(BreakerState {
198 consecutive_failures: 0,
199 open_until: None,
200 last_failure_at: now,
201 trip_count: 0,
202 });
203 state.last_failure_at = now;
204 state.consecutive_failures = state.consecutive_failures.saturating_add(1);
205 if state.consecutive_failures < MAX_CONSECUTIVE_SUOTAR_FAILURES {
206 return None;
207 }
208 state.trip_count = state.trip_count.saturating_add(1);
209 let cooldown = base_cooldown * state.trip_count.min(MAX_COOLDOWN_TRIPS);
210 state.open_until = Some(now + cooldown);
211 Some(BreakerTrip {
212 cooldown,
213 consecutive_failures: state.consecutive_failures,
214 trip_count: state.trip_count,
215 })
216}
217
218#[cfg(test)]
219fn reset(scope: &ScopeKey) {
220 let mut breakers = BREAKERS.lock();
221 breakers.retain(|(key_scope, _), _| key_scope != scope);
222}
223
224#[cfg(test)]
225mod tests {
226 use uuid::Uuid;
227
228 use super::*;
229
230 const TARGET: BreakerTarget = BreakerTarget::StudyRegistry;
231
232 fn key() -> ScopeKey {
233 ScopeKey::Course(Uuid::new_v4())
234 }
235
236 #[test]
237 fn the_breaker_opens_only_after_the_documented_run_of_failures() {
238 let key = key();
239 for _ in 1..MAX_CONSECUTIVE_SUOTAR_FAILURES {
240 assert!(record_failure(&key, TARGET, cooldown(false)).is_none());
241 assert!(!is_open(&key, TARGET));
242 }
243 assert!(record_failure(&key, TARGET, cooldown(false)).is_some());
244 assert!(is_open(&key, TARGET));
245 reset(&key);
246 }
247
248 #[test]
249 fn one_success_puts_the_run_of_failures_back_to_zero() {
250 let key = key();
251 for _ in 1..MAX_CONSECUTIVE_SUOTAR_FAILURES {
252 record_failure(&key, TARGET, cooldown(false));
253 }
254 record_success(&key, TARGET);
255 assert!(record_failure(&key, TARGET, cooldown(false)).is_none());
256 assert!(!is_open(&key, TARGET));
257 reset(&key);
258 }
259
260 #[test]
261 fn two_scopes_do_not_trip_each_other() {
262 let storm = key();
263 let bystander = key();
264 for _ in 0..MAX_CONSECUTIVE_SUOTAR_FAILURES {
265 record_failure(&storm, TARGET, cooldown(false));
266 }
267 assert!(is_open(&storm, TARGET));
268 assert!(!is_open(&bystander, TARGET));
269 reset(&storm);
270 reset(&bystander);
271 }
272
273 #[test]
274 fn a_tripped_breaker_closes_once_its_cooldown_has_elapsed() {
275 let key = key();
276 for _ in 0..MAX_CONSECUTIVE_SUOTAR_FAILURES {
277 record_failure(&key, TARGET, Duration::ZERO);
278 }
279 assert!(!is_open(&key, TARGET));
280 reset(&key);
281 }
282}