Skip to main content

headless_lms_credit_registration/runtime/
process_local.rs

1//! The in-memory maps this worker process keeps its advisory state in: the circuit breakers', the
2//! rate limiter's, and the worker loop's last logged skips.
3
4use std::collections::HashMap;
5use std::hash::Hash;
6use std::sync::{LazyLock, Mutex, MutexGuard};
7use std::time::{Duration, Instant};
8
9use headless_lms_models::credit_registrations::RegistrationScope;
10use uuid::Uuid;
11
12/// Whose share of the breakers and the limiter a run uses. Keyed by scope rather than global: a
13/// test driving a deliberate outage for its own course must not silence the pipeline for every
14/// other test running at the same moment. Production only ever uses the global key.
15#[derive(Debug, Clone, PartialEq, Eq, Hash)]
16pub(super) enum ScopeKey {
17    /// Production, and any unscoped run.
18    Global,
19    Course(Uuid),
20    User(Uuid),
21    Registrations(Vec<Uuid>),
22}
23
24impl ScopeKey {
25    pub(super) fn of(scope: &RegistrationScope) -> Self {
26        if let Some(course_id) = scope.course_id {
27            Self::Course(course_id)
28        } else if let Some(user_id) = scope.user_id {
29            Self::User(user_id)
30        } else if !scope.credit_registration_ids.is_empty() {
31            let mut ids = scope.credit_registration_ids.clone();
32            ids.sort();
33            Self::Registrations(ids)
34        } else {
35            Self::Global
36        }
37    }
38}
39
40/// A map private to this worker process, for state that is advisory: a poisoned lock is recovered
41/// rather than taking the worker down.
42pub(super) struct ProcessLocalMap<K, V>(LazyLock<Mutex<HashMap<K, V>>>);
43
44impl<K, V> ProcessLocalMap<K, V> {
45    pub(super) const fn new() -> Self {
46        Self(LazyLock::new(|| Mutex::new(HashMap::new())))
47    }
48
49    pub(super) fn lock(&self) -> MutexGuard<'_, HashMap<K, V>> {
50        self.0.lock().unwrap_or_else(|poisoned| {
51            warn!("Recovered a poisoned breaker/rate-limit lock after a panic elsewhere");
52            poisoned.into_inner()
53        })
54    }
55}
56
57/// How often an unchanged report is written anyway, so its `updated_at` still shows the process is
58/// alive.
59const UNCHANGED_REPORT_REFRESH: Duration = Duration::from_secs(60);
60
61/// What this process last reported per key, to the dashboard's copy of its in-memory state or to
62/// the log.
63pub(super) struct LastReported<K, S> {
64    reports: ProcessLocalMap<K, (S, Instant)>,
65    /// How often an unchanged report is due anyway; `None` for never.
66    refresh: Option<Duration>,
67}
68
69impl<K: Eq + Hash, S> LastReported<K, S> {
70    /// Refreshed every [`UNCHANGED_REPORT_REFRESH`].
71    pub(super) const fn new() -> Self {
72        Self {
73            reports: ProcessLocalMap::new(),
74            refresh: Some(UNCHANGED_REPORT_REFRESH),
75        }
76    }
77
78    /// Due again only once the report changes or [`Self::clear`] forgets it.
79    pub(super) const fn without_refresh() -> Self {
80        Self {
81            reports: ProcessLocalMap::new(),
82            refresh: None,
83        }
84    }
85
86    /// Whether `report` differs from the last one recorded for `key` by `is_same`, or that one has
87    /// aged past the refresh interval.
88    pub(super) fn is_due(&self, key: &K, report: &S, is_same: impl FnOnce(&S, &S) -> bool) -> bool {
89        self.reports.lock().get(key).is_none_or(|(written, at)| {
90            !is_same(written, report) || self.refresh.is_some_and(|refresh| at.elapsed() >= refresh)
91        })
92    }
93
94    /// Notes that `report` was written for `key`; only after the write succeeded.
95    pub(super) fn record(&self, key: K, report: S) {
96        self.reports.lock().insert(key, (report, Instant::now()));
97    }
98
99    /// Forgets `key`'s last report, so the next one is due whatever it is.
100    pub(super) fn clear(&self, key: &K) {
101        self.reports.lock().remove(key);
102    }
103}
104
105#[cfg(test)]
106mod tests {
107    use super::*;
108
109    #[test]
110    fn a_scoped_run_gets_its_own_key_and_an_unscoped_one_gets_the_global_key() {
111        let course = Uuid::new_v4();
112        let user = Uuid::new_v4();
113        assert_eq!(
114            ScopeKey::of(&RegistrationScope::default()),
115            ScopeKey::Global
116        );
117        assert_eq!(
118            ScopeKey::of(&RegistrationScope::for_course(course)),
119            ScopeKey::Course(course)
120        );
121        assert_eq!(
122            ScopeKey::of(&RegistrationScope {
123                user_id: Some(user),
124                ..RegistrationScope::default()
125            }),
126            ScopeKey::User(user)
127        );
128    }
129
130    #[test]
131    fn a_registration_scope_is_order_independent() {
132        let first = Uuid::new_v4();
133        let second = Uuid::new_v4();
134        let one = RegistrationScope {
135            credit_registration_ids: vec![first, second],
136            ..RegistrationScope::default()
137        };
138        let other = RegistrationScope {
139            credit_registration_ids: vec![second, first],
140            ..RegistrationScope::default()
141        };
142        assert_eq!(ScopeKey::of(&one), ScopeKey::of(&other));
143    }
144}