headless_lms_credit_registration/runtime/suotar/
health_report.rs1use 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
12pub(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
48pub(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}