Skip to main content

headless_lms_server/controllers/main_frontend/credit_registration_admin/
phases.rs

1//! The Workers tab: one row per pipeline phase, with the queue each is responsible for.
2//!
3//! Phases, not pods. Pausing `import` while `verify` keeps running is a real incident move that no
4//! pod-level control can express, and the pause here is our own flag — the k8s status page still
5//! answers whether the process hosting a phase is up.
6
7use 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/// One phase as the Workers tab renders it.
28///
29/// Wider than `CreditRegistrationPhaseStatus`, which the pause/resume/run-now responses return: this
30/// one also carries the last error, the run window and the queue.
31#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
32pub struct CreditRegistrationPhaseRow {
33    pub phase: String,
34    /// The worker process whose loop runs the phase. Rows are grouped by it, because a dead pod
35    /// makes every phase inside it go stale at once and that reads as one fault, not seven.
36    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    /// Our own wording or the study registry's code, never its prose.
47    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    /// False only for a phase-state row whose name is no `CreditRegistrationPhase`, which no worker
52    /// runs or reports for.
53    pub is_known_phase: bool,
54    /// Computed server-side: a page comparing its own clock against a server timestamp misjudges
55    /// this on a skewed client.
56    pub seconds_since_heartbeat: Option<i64>,
57    pub last_run_duration_secs: Option<i64>,
58    /// Always `false` while paused or never heartbeated.
59    pub heartbeat_late: bool,
60    /// The same verdict the `PhaseFailing` alert reaches; see `is_phase_failing`.
61    pub failing: bool,
62    /// The ledger states nothing but this phase moves a row out of. Empty for the phases whose work
63    /// is not a ledger state: `materialize` waits on completions, the syncer's phases on modules.
64    pub owned_states: Vec<CreditRegistrationState>,
65    /// Live rows in `owned_states` waiting on this phase, of `no_usable_enrolment` only those due a
66    /// check, or `None` where there are none to own — which is not the same as an empty queue.
67    pub queue_depth: Option<i64>,
68}
69
70/// Where one circuit breaker stands, as its worker last reported it.
71#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone, Copy, ToSchema)]
72#[serde(rename_all = "snake_case")]
73pub enum CircuitBreakerStatus {
74    Closed,
75    /// In its cooldown: the phases it pauses skip their iterations.
76    Open,
77    /// Past its cooldown, and the next iteration sends a single-item probe that closes it only if
78    /// it succeeds.
79    WaitingToProbe,
80}
81
82/// One worker process's circuit breaker, as the worker last reported it.
83#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
84pub struct CreditRegistrationCircuitBreakerState {
85    pub process_name: String,
86    pub target: BreakerTarget,
87    /// The endpoints whose phases the breaker pauses.
88    pub endpoints: Vec<SuotarEndpoint>,
89    pub status: CircuitBreakerStatus,
90    pub consecutive_failures: i64,
91    /// How much of the cooldown is left. Computed server-side, like `seconds_since_heartbeat`.
92    pub open_for_secs: Option<i64>,
93    pub trip_count: i64,
94    /// When the worker last reported the state.
95    pub updated_at: DateTime<Utc>,
96}
97
98#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
99pub struct CreditRegistrationPhaseList {
100    /// In process then pipeline order, so the grouping is a fold over the list.
101    pub phases: Vec<CreditRegistrationPhaseRow>,
102    pub heartbeat_interval_multiplier: i32,
103    pub consecutive_failure_limit: i32,
104    /// Every phase is stopped, which is what the kill switch does.
105    pub paused_globally: bool,
106    pub circuit_breakers: Vec<CreditRegistrationCircuitBreakerState>,
107}
108
109/**
110GET `/api/v0/main-frontend/credit-registration-admin/phases` - Every pipeline phase, its heartbeat
111and the queue it is responsible for, and the workers' circuit breakers.
112*/
113#[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}