Skip to main content

headless_lms_credit_registration/runtime/suotar/
breaker.rs

1//! The circuit breakers the study-registry phases share within one worker process: one every such
2//! phase stops for, and one only the phase that submits to Sisu stops for, since Sisu timing out on
3//! submissions says nothing about the rest of Suotar.
4//!
5//! `BREAKERS` is a process-local static: `credit-registrar` and `suotar-syncer` are separate OS
6//! processes (see [`crate::WorkerProcess`]), each with its own map, so an outage tripping the
7//! breaker in one does not pause the study-registry phases of the other. Only the phases within the
8//! same process actually share a breaker per scope key.
9//!
10//! Keyed by [`ScopeKey`], so a test driving a deliberate outage for its own course does not silence
11//! the pipeline for every other test running at the same moment.
12
13use 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;
22/// The first cooldown; each trip without a success between adds another, up to
23/// [`MAX_COOLDOWN_TRIPS`] of them.
24const SUOTAR_COOLDOWN: Duration = Duration::from_secs(5 * 60);
25const MAX_COOLDOWN_TRIPS: u32 = 3;
26/// Playwright's per-test budget is 100 s, which the production cooldown does not fit inside: a test
27/// that trips the breaker deliberately has to be able to watch it recover.
28const TEST_SUOTAR_COOLDOWN: Duration = Duration::from_secs(5);
29
30pub(super) use headless_lms_models::suotar_circuit_breakers::BreakerTarget;
31
32/// How long a run of failures that never tripped the breaker is remembered, so a scope never run
33/// again leaves the map. Failures an outage spreads between hour-long timed-out calls must still
34/// add up.
35const 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    /// Times opened since the last success, which the next cooldown grows with.
49    trip_count: u32,
50}
51
52impl BreakerState {
53    /// Whether the entry still says anything: an open cooldown, or a recent enough run of failures.
54    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
64/// The global breakers as this process last wrote them for the dashboard.
65pub(super) static REPORTED: LastReported<BreakerTarget, BreakerSnapshot> = LastReported::new();
66
67/// The first cooldown, which later trips multiply.
68pub(super) fn cooldown(test_mode: bool) -> Duration {
69    if test_mode {
70        TEST_SUOTAR_COOLDOWN
71    } else {
72        SUOTAR_COOLDOWN
73    }
74}
75
76/// Whether the phases `target` covers should skip this iteration.
77pub(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    // Dropped rather than reset in place so an idle scope leaves the map; the fresh entry the next
91    // failure creates is the state a reset would have left behind anyway.
92    breakers.remove(&key);
93    false
94}
95
96/// Whether the breaker's cooldown has ended with no success since: the next iteration is a probe,
97/// and sends one item only.
98pub(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
109/// Whether a breaker with this run of failures and cooldown is past the cooldown without the
110/// success that closes it: what `is_half_open` asks of a live breaker, for one read back from the
111/// database.
112pub 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/// What one breaker holds right now, in this process.
122#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
123pub(super) struct BreakerSnapshot {
124    pub open: bool,
125    pub consecutive_failures: u32,
126    /// When the cooldown ends, or ended for a breaker that is half-open now.
127    pub open_until: Option<DateTime<Utc>>,
128    pub trip_count: u32,
129}
130
131impl BreakerSnapshot {
132    /// Whether the two write the same database row, `open_until` to the second: converting it from
133    /// the monotonic clock moves it a little on every [`snapshot`].
134    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
142/// Reads a breaker without touching it, for the dashboard. Not [`is_open`], which clears an elapsed
143/// cooldown as a side effect.
144pub(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
169/// Returns the number of times this breaker had tripped, if this success closed it.
170pub(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
178/// What one failure that tripped or re-tripped a breaker did, for the caller's log line.
179pub(super) struct BreakerTrip {
180    pub cooldown: Duration,
181    pub consecutive_failures: u32,
182    pub trip_count: u32,
183}
184
185/// Returns what this failure did to the breaker, if it tripped or re-tripped it. `base_cooldown` is
186/// the first trip's; a failure while half-open trips again at once, for longer.
187pub(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}