headless_lms_server/controllers/main_frontend/credit_registration_admin/
phases.rs1use std::collections::HashMap;
8
9use headless_lms_models::credit_registration_phase_state::{
10 self, CreditRegistrationPhaseState as PhaseStateRow,
11};
12use headless_lms_models::credit_registrations::{self, CreditRegistrationState};
13use headless_lms_models::suotar_api_calls::SuotarEndpoint;
14use headless_lms_models::suotar_circuit_breakers::{self, BreakerTarget, SuotarCircuitBreaker};
15use utoipa::ToSchema;
16
17use crate::domain::credit_registration::health::{
18 PHASE_CONSECUTIVE_FAILURE_LIMIT, PHASE_HEARTBEAT_INTERVAL_MULTIPLIER, is_heartbeat_late,
19 is_phase_failing,
20};
21use crate::prelude::*;
22use headless_lms_credit_registration::CreditRegistrationPhase;
23use headless_lms_credit_registration::registry_health::{endpoints_paused_by, is_waiting_to_probe};
24
25use super::authorize_credit_registration_admin;
26
27#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
32pub struct CreditRegistrationPhaseRow {
33 pub phase: String,
34 pub process_name: String,
37 pub expected_interval_secs: i32,
38 pub last_heartbeat_at: Option<DateTime<Utc>>,
39 pub last_run_started_at: Option<DateTime<Utc>>,
40 pub last_run_finished_at: Option<DateTime<Utc>>,
41 pub last_success_at: Option<DateTime<Utc>>,
42 pub next_run_at: Option<DateTime<Utc>>,
43 pub items_processed_last_run: Option<i32>,
44 pub items_failed_last_run: Option<i32>,
45 pub consecutive_failures: i32,
46 pub last_error: Option<String>,
48 pub paused_at: Option<DateTime<Utc>>,
49 pub paused_by_user_id: Option<Uuid>,
50 pub pause_reason: Option<String>,
51 pub is_known_phase: bool,
54 pub seconds_since_heartbeat: Option<i64>,
57 pub last_run_duration_secs: Option<i64>,
58 pub heartbeat_late: bool,
60 pub failing: bool,
62 pub owned_states: Vec<CreditRegistrationState>,
65 pub queue_depth: Option<i64>,
68}
69
70#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, Copy, ToSchema)]
72#[serde(rename_all = "snake_case")]
73pub enum CircuitBreakerStatus {
74 Closed,
75 Open,
77 WaitingToProbe,
80}
81
82#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
84pub struct CreditRegistrationCircuitBreakerState {
85 pub process_name: String,
86 pub target: BreakerTarget,
87 pub endpoints: Vec<SuotarEndpoint>,
89 pub status: CircuitBreakerStatus,
90 pub consecutive_failures: i64,
91 pub open_for_secs: Option<i64>,
93 pub trip_count: i64,
94 pub updated_at: DateTime<Utc>,
96}
97
98#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
99pub struct CreditRegistrationPhaseList {
100 pub phases: Vec<CreditRegistrationPhaseRow>,
102 pub heartbeat_interval_multiplier: i32,
103 pub consecutive_failure_limit: i32,
104 pub paused_globally: bool,
106 pub circuit_breakers: Vec<CreditRegistrationCircuitBreakerState>,
107}
108
109#[instrument(skip(pool))]
114#[utoipa::path(
115 get,
116 path = "/phases",
117 operation_id = "listCreditRegistrationPhases",
118 tag = "credit-registration-admin",
119 responses(
120 (status = 200, description = "One row per pipeline phase, and the workers' circuit breakers", body = CreditRegistrationPhaseList)
121 )
122)]
123pub async fn list_credit_registration_phases(
124 user: AuthUser,
125 pool: web::Data<PgPool>,
126) -> ControllerResult<web::Json<CreditRegistrationPhaseList>> {
127 let mut conn = pool.acquire().await?;
128 let token = authorize_credit_registration_admin(&mut conn, user.id).await?;
129
130 let depths: HashMap<CreditRegistrationState, i64> =
131 credit_registrations::count_by_state(&mut conn)
132 .await?
133 .into_iter()
134 .collect();
135 let due_enrolment_checks = credit_registrations::count_due_enrolment_checks(&mut conn).await?;
136 let now = Utc::now();
137 let mut phases: Vec<CreditRegistrationPhaseRow> =
138 credit_registration_phase_state::get_all(&mut conn)
139 .await?
140 .into_iter()
141 .map(|row| to_phase_row(row, now, &depths, due_enrolment_checks))
142 .collect();
143 phases.sort_by_key(|row| {
144 (
145 row.process_name.clone(),
146 CreditRegistrationPhase::from_phase_name(&row.phase)
147 .map_or(usize::MAX, CreditRegistrationPhase::pipeline_index),
148 )
149 });
150 let circuit_breakers = suotar_circuit_breakers::get_all(&mut conn)
151 .await?
152 .into_iter()
153 .map(|breaker| to_circuit_breaker_state(breaker, now))
154 .collect();
155
156 token.authorized_ok(web::Json(CreditRegistrationPhaseList {
157 paused_globally: !phases.is_empty() && phases.iter().all(|row| row.paused_at.is_some()),
158 phases,
159 heartbeat_interval_multiplier: PHASE_HEARTBEAT_INTERVAL_MULTIPLIER,
160 consecutive_failure_limit: PHASE_CONSECUTIVE_FAILURE_LIMIT,
161 circuit_breakers,
162 }))
163}
164
165fn to_phase_row(
166 row: PhaseStateRow,
167 now: DateTime<Utc>,
168 depths: &HashMap<CreditRegistrationState, i64>,
169 due_enrolment_checks: i64,
170) -> CreditRegistrationPhaseRow {
171 let known = CreditRegistrationPhase::from_phase_name(&row.phase);
172 let owned_states: Vec<CreditRegistrationState> = known
173 .map(|phase| phase.owned_states().to_vec())
174 .unwrap_or_default();
175 let seconds_since_heartbeat = row.last_heartbeat_at.map(|at| (now - at).num_seconds());
176 let heartbeat_late = is_heartbeat_late(
177 row.last_heartbeat_at,
178 row.expected_interval_secs,
179 row.paused_at,
180 now,
181 );
182 let depth_of = |state| depths.get(&state).copied().unwrap_or(0);
183 let failing = is_phase_failing(&row, now, depth_of, due_enrolment_checks);
184 CreditRegistrationPhaseRow {
185 is_known_phase: known.is_some(),
186 queue_depth: known
187 .filter(|_| !owned_states.is_empty())
188 .map(|phase| phase.queue_depth(depth_of, due_enrolment_checks)),
189 owned_states,
190 seconds_since_heartbeat,
191 heartbeat_late,
192 failing,
193 last_run_duration_secs: row
194 .last_run_started_at
195 .zip(row.last_run_finished_at)
196 .map(|(started, finished)| (finished - started).num_seconds()),
197 phase: row.phase,
198 process_name: row.process_name,
199 expected_interval_secs: row.expected_interval_secs,
200 last_heartbeat_at: row.last_heartbeat_at,
201 last_run_started_at: row.last_run_started_at,
202 last_run_finished_at: row.last_run_finished_at,
203 last_success_at: row.last_success_at,
204 next_run_at: row.next_run_at,
205 items_processed_last_run: row.items_processed_last_run,
206 items_failed_last_run: row.items_failed_last_run,
207 consecutive_failures: row.consecutive_failures,
208 last_error: row.last_error,
209 paused_at: row.paused_at,
210 paused_by_user_id: row.paused_by_user_id,
211 pause_reason: row.pause_reason,
212 }
213}
214
215pub fn _add_routes(cfg: &mut ServiceConfig) {
216 cfg.route("/phases", web::get().to(list_credit_registration_phases));
217}
218
219fn to_circuit_breaker_state(
220 breaker: SuotarCircuitBreaker,
221 now: DateTime<Utc>,
222) -> CreditRegistrationCircuitBreakerState {
223 let open_for_secs = breaker
224 .open_until
225 .map(|until| (until - now).num_seconds())
226 .filter(|&secs| secs > 0);
227 let consecutive_failures = u32::try_from(breaker.consecutive_failures).unwrap_or_default();
228 let status = if breaker.open_until.is_some_and(|until| now < until) {
229 CircuitBreakerStatus::Open
230 } else if is_waiting_to_probe(consecutive_failures, breaker.open_until, now) {
231 CircuitBreakerStatus::WaitingToProbe
232 } else {
233 CircuitBreakerStatus::Closed
234 };
235 let endpoints = endpoints_paused_by(&breaker.process_name, breaker.target);
236 CreditRegistrationCircuitBreakerState {
237 process_name: breaker.process_name,
238 target: breaker.target,
239 endpoints,
240 status,
241 consecutive_failures: i64::from(breaker.consecutive_failures),
242 open_for_secs,
243 trip_count: i64::from(breaker.trip_count),
244 updated_at: breaker.updated_at,
245 }
246}