Skip to main content

headless_lms_credit_registration/runtime/suotar/
health_report.rs

1//! Copying this process's breaker and limiter state to the database for the dashboard, which runs
2//! in another process.
3
4use headless_lms_models::{suotar_circuit_breakers, suotar_endpoint_rate_limits};
5use sqlx::PgConnection;
6
7use super::{breaker, rate_limit};
8use crate::error::CreditRegistrationResult;
9use crate::phase::CreditRegistrationPhase;
10use crate::runtime::process_local::ScopeKey;
11
12/// Copies the state of the breakers that pause the phase to the database for the dashboard, which
13/// runs in another process.
14pub(in crate::runtime) async fn report_breakers(
15    conn: &mut PgConnection,
16    phase: CreditRegistrationPhase,
17) -> CreditRegistrationResult<()> {
18    let spec = phase.spec();
19    for &target in spec.breakers {
20        let breaker = breaker::snapshot(&ScopeKey::Global, target);
21        if !breaker::REPORTED.is_due(&target, &breaker, breaker::BreakerSnapshot::is_same_report) {
22            continue;
23        }
24        trace!(
25            ?target,
26            consecutive_failures = breaker.consecutive_failures,
27            open = breaker.open,
28            trip_count = breaker.trip_count,
29            "Reporting circuit breaker state to the dashboard"
30        );
31        suotar_circuit_breakers::upsert(
32            conn,
33            &suotar_circuit_breakers::SuotarCircuitBreakerReport {
34                process_name: spec.process.as_str(),
35                target,
36                consecutive_failures: i32::try_from(breaker.consecutive_failures)
37                    .unwrap_or(i32::MAX),
38                open_until: breaker.open_until,
39                trip_count: i32::try_from(breaker.trip_count).unwrap_or(i32::MAX),
40            },
41        )
42        .await?;
43        breaker::REPORTED.record(target, breaker);
44    }
45    Ok(())
46}
47
48/// Copies the limiter state of the phase's endpoints to the database for the dashboard, which runs
49/// in another process.
50pub(in crate::runtime) async fn report_rate_limits(
51    conn: &mut PgConnection,
52    phase: CreditRegistrationPhase,
53) -> CreditRegistrationResult<()> {
54    for endpoint in super::study_registry_endpoints(phase) {
55        let Some(limiter) = rate_limit::snapshot(&ScopeKey::Global, endpoint) else {
56            continue;
57        };
58        if !rate_limit::REPORTED.is_due(&endpoint, &limiter, PartialEq::eq) {
59            continue;
60        }
61        trace!(
62            ?endpoint,
63            rate_share = limiter.share,
64            available = limiter.available,
65            "Reporting rate limit state to the dashboard"
66        );
67        suotar_endpoint_rate_limits::upsert(
68            conn,
69            &suotar_endpoint_rate_limits::SuotarEndpointRateLimitReport {
70                endpoint,
71                rate_share: limiter.share as f32,
72                full_rate_per_minute: limiter.rate.per_minute as i32,
73                available: i32::try_from(limiter.available).unwrap_or(i32::MAX),
74            },
75        )
76        .await?;
77        rate_limit::REPORTED.record(endpoint, limiter);
78    }
79    Ok(())
80}