Skip to main content

headless_lms_server/controllers/main_frontend/credit_registration_admin/
dashboard.rs

1//! The Overview tab, the Suotar health panel and phase pause/resume/run-now controls.
2
3use std::collections::HashMap;
4
5use headless_lms_models::credit_registration_admin_actions::{
6    CreditRegistrationAdminAction, CreditRegistrationAdminActionTarget, GLOBAL_ADMIN_ROLE,
7    NewCreditRegistrationAdminAction,
8};
9use headless_lms_models::credit_registration_phase_state;
10use headless_lms_models::credit_registrations::{
11    self, CreditRegistrationErrorCode, CreditRegistrationErrorCodeCount, CreditRegistrationState,
12    OldestNonTerminalRegistration, StuckRegistrationCount,
13};
14use headless_lms_models::library::credit_registration::PendingReasonCounts;
15use headless_lms_models::suotar_api_calls::{
16    self, SuotarEndpoint, SuotarEndpointStanding as SuotarEndpointStandingRow,
17    SuotarEndpointStatsForWindow,
18};
19use utoipa::ToSchema;
20
21use crate::domain::credit_registration::health::{
22    CreditRegistrationHealth, evaluate, is_heartbeat_late, stuck_thresholds,
23};
24use crate::prelude::*;
25use headless_lms_credit_registration::CreditRegistrationPhase;
26
27use super::{ATTENTION_TOO_MANY_ATTEMPTS, authorize_credit_registration_admin, required_reason};
28
29const THROUGHPUT_DAYS: i64 = 30;
30
31const ENDPOINT_STATS_WINDOWS_SECS: [i64; 3] = [60 * 60, 24 * 60 * 60, 7 * 24 * 60 * 60];
32
33#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
34pub struct CreditRegistrationStateTotal {
35    pub state: CreditRegistrationState,
36    pub count: i64,
37}
38
39#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
40pub struct CreditRegistrationErrorCodeTotal {
41    pub error_code: CreditRegistrationErrorCode,
42    /// Rows the pipeline is still working on.
43    pub in_flight_count: i64,
44    /// Rows that ended on this code.
45    pub terminal_failure_count: i64,
46}
47
48#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
49pub struct CreditRegistrationOldestNonTerminal {
50    pub credit_registration_id: Uuid,
51    pub state: CreditRegistrationState,
52    pub state_entered_at: DateTime<Utc>,
53    /// Computed server-side: a page comparing its own clock against a server timestamp misjudges
54    /// this on a skewed client, the same reason `seconds_since_heartbeat` is computed here too.
55    pub seconds_in_state: i64,
56}
57
58#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
59pub struct CreditRegistrationThroughputBucket {
60    pub day: DateTime<Utc>,
61    pub registered_count: i64,
62    /// `duplicate` and `not_improved`: the credit exists, and we did not put it there.
63    pub other_success_count: i64,
64    pub failed_count: i64,
65}
66
67#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
68pub struct CreditRegistrationStuckTotal {
69    pub state: CreditRegistrationState,
70    pub count: i64,
71    pub severely_stuck_count: i64,
72    pub oldest_state_entered_at: Option<DateTime<Utc>>,
73}
74
75/// Where one study registry endpoint stands, over all time.
76#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
77pub struct SuotarEndpointStanding {
78    pub endpoint: SuotarEndpoint,
79    pub last_success_at: Option<DateTime<Utc>>,
80    pub last_failure_at: Option<DateTime<Utc>>,
81    pub consecutive_failures: i64,
82}
83
84/// One pipeline phase's heartbeat, written by the worker loops and by unscoped runs only, never by a
85/// narrowed one. Returned by the pause/resume/run-now actions; the Workers tab lists
86/// `CreditRegistrationPhaseRow` instead, which is wider.
87#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
88pub struct CreditRegistrationPhaseStatus {
89    pub phase: String,
90    pub process_name: String,
91    pub expected_interval_secs: i32,
92    pub last_heartbeat_at: Option<DateTime<Utc>>,
93    pub last_success_at: Option<DateTime<Utc>>,
94    pub last_run_finished_at: Option<DateTime<Utc>>,
95    pub items_processed_last_run: Option<i32>,
96    pub items_failed_last_run: Option<i32>,
97    pub consecutive_failures: i32,
98    pub paused_at: Option<DateTime<Utc>>,
99    pub pause_reason: Option<String>,
100    /// False only for a phase-state row whose name is no `CreditRegistrationPhase`, which no worker
101    /// runs or reports for.
102    pub is_known_phase: bool,
103    /// Computed server-side: a page comparing its own clock against a server timestamp misjudges this
104    /// on a skewed client.
105    pub seconds_since_heartbeat: Option<i64>,
106    /// `seconds_since_heartbeat > expected_interval_secs * health.thresholds.phase_heartbeat_interval_multiplier`.
107    /// Always `false` while paused or never heartbeated.
108    pub heartbeat_late: bool,
109}
110
111#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
112pub struct CreditRegistrationOverview {
113    pub health: CreditRegistrationHealth,
114    pub counts_by_state: Vec<CreditRegistrationStateTotal>,
115    /// The `pending` depth split by what each row is waiting on, which the ledger does not store.
116    pub pending_by_reason: PendingReasonCounts,
117    pub error_codes: Vec<CreditRegistrationErrorCodeTotal>,
118    /// Live rows a detector picked or the pipeline flagged. The one definition of "needs a human":
119    /// `/attention` pages through exactly these rows and reports the same total.
120    pub needs_admin_attention_count: i64,
121    pub oldest_non_terminal: Option<CreditRegistrationOldestNonTerminal>,
122    pub throughput: Vec<CreditRegistrationThroughputBucket>,
123    pub throughput_days: i64,
124    pub stuck: Vec<CreditRegistrationStuckTotal>,
125    pub endpoints: Vec<SuotarEndpointStanding>,
126}
127
128#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
129pub struct SuotarEndpointWindowStats {
130    pub endpoint: SuotarEndpoint,
131    pub call_count: i64,
132    pub failed_call_count: i64,
133    pub in_flight_count: i64,
134    pub ok_item_count: i64,
135    pub error_item_count: i64,
136    pub pending_item_count: i64,
137    pub p50_duration_ms: Option<i32>,
138    pub p95_duration_ms: Option<i32>,
139    pub last_success_at: Option<DateTime<Utc>>,
140    pub last_failure_at: Option<DateTime<Utc>>,
141    /// The registry's own request-level code, an identifier rather than prose.
142    pub last_request_level_error_code: Option<String>,
143}
144
145#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
146pub struct SuotarHealthWindow {
147    pub window_secs: i64,
148    pub endpoints: Vec<SuotarEndpointWindowStats>,
149}
150
151#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
152pub struct SuotarHealth {
153    pub windows: Vec<SuotarHealthWindow>,
154}
155
156#[derive(Debug, Deserialize, ToSchema)]
157pub struct AdminPausePhasePayload {
158    pub reason: String,
159}
160
161#[derive(Debug, Deserialize, ToSchema)]
162pub struct AdminPhaseActionPayload {
163    pub reason: Option<String>,
164}
165
166/**
167GET `/api/v0/main-frontend/credit-registration-admin/overview` - Everything the Overview tab and the
168alert banner render, in one request so the tiles cannot contradict each other.
169*/
170#[instrument(skip(pool))]
171#[utoipa::path(
172    get,
173    path = "/overview",
174    operation_id = "getCreditRegistrationOverview",
175    tag = "credit-registration-admin",
176    responses(
177        (status = 200, description = "Counts, throughput, phase heartbeats and the active alerts", body = CreditRegistrationOverview)
178    )
179)]
180pub async fn get_credit_registration_overview(
181    user: AuthUser,
182    pool: web::Data<PgPool>,
183) -> ControllerResult<web::Json<CreditRegistrationOverview>> {
184    let mut conn = pool.acquire().await?;
185    let token = authorize_credit_registration_admin(&mut conn, user.id).await?;
186
187    let stuck_rows = credit_registrations::count_stuck(&mut conn, &stuck_thresholds()).await?;
188    let depths = credit_registrations::count_by_state(&mut conn).await?;
189    let health = evaluate(&mut conn, &stuck_rows, &depths).await?;
190    let counts_by_state = depths
191        .iter()
192        .map(|&(state, count)| CreditRegistrationStateTotal { state, count })
193        .collect();
194    let pending_by_reason = credit_registrations::count_pending_by_reason(&mut conn).await?;
195    let error_codes = credit_registrations::count_by_error_code(&mut conn)
196        .await?
197        .into_iter()
198        .map(to_error_code_total)
199        .collect();
200    let needs_admin_attention_count = credit_registrations::count_needing_attention(
201        &mut conn,
202        &stuck_thresholds(),
203        ATTENTION_TOO_MANY_ATTEMPTS,
204    )
205    .await?
206    .map_or(0, |row| row.total_count);
207    let oldest_non_terminal = credit_registrations::get_oldest_non_terminal(&mut conn)
208        .await?
209        .map(|row| to_oldest_non_terminal(row, Utc::now()));
210    let throughput = credit_registrations::get_throughput_by_day(
211        &mut conn,
212        Utc::now() - chrono::Duration::days(THROUGHPUT_DAYS),
213    )
214    .await?
215    .into_iter()
216    .map(|row| CreditRegistrationThroughputBucket {
217        day: row.day,
218        registered_count: row.registered_count,
219        other_success_count: row.other_success_count,
220        failed_count: row.failed_count,
221    })
222    .collect();
223    let stuck = stuck_rows.into_iter().map(to_stuck_total).collect();
224    let endpoints = suotar_api_calls::get_endpoint_standings(&mut conn)
225        .await?
226        .into_iter()
227        .map(to_endpoint_standing)
228        .collect();
229
230    token.authorized_ok(web::Json(CreditRegistrationOverview {
231        health,
232        counts_by_state,
233        pending_by_reason,
234        error_codes,
235        needs_admin_attention_count,
236        oldest_non_terminal,
237        throughput,
238        throughput_days: THROUGHPUT_DAYS,
239        stuck,
240        endpoints,
241    }))
242}
243
244/**
245GET `/api/v0/main-frontend/credit-registration-admin/suotar-health` - Per-endpoint call counts,
246success rates and latency percentiles over an hour, a day and a week.
247*/
248#[instrument(skip(pool))]
249#[utoipa::path(
250    get,
251    path = "/suotar-health",
252    operation_id = "getSuotarHealth",
253    tag = "credit-registration-admin",
254    responses(
255        (status = 200, description = "Study registry traffic per endpoint and window", body = SuotarHealth)
256    )
257)]
258pub async fn get_suotar_health(
259    user: AuthUser,
260    pool: web::Data<PgPool>,
261) -> ControllerResult<web::Json<SuotarHealth>> {
262    let mut conn = pool.acquire().await?;
263    let token = authorize_credit_registration_admin(&mut conn, user.id).await?;
264
265    let mut by_window: HashMap<i64, Vec<SuotarEndpointWindowStats>> = HashMap::new();
266    for row in
267        suotar_api_calls::get_endpoint_stats_for_windows(&mut conn, &ENDPOINT_STATS_WINDOWS_SECS)
268            .await?
269    {
270        by_window
271            .entry(row.window_secs)
272            .or_default()
273            .push(to_endpoint_window_stats_for_window(row));
274    }
275    let windows = ENDPOINT_STATS_WINDOWS_SECS
276        .into_iter()
277        .map(|window_secs| SuotarHealthWindow {
278            window_secs,
279            endpoints: by_window.remove(&window_secs).unwrap_or_default(),
280        })
281        .collect();
282
283    token.authorized_ok(web::Json(SuotarHealth { windows }))
284}
285
286/**
287POST `/api/v0/main-frontend/credit-registration-admin/phases/{phase}/pause` - Pauses one phase: the
288worker loop skips it on every tick until it is resumed.
289*/
290#[instrument(skip(pool, payload))]
291#[utoipa::path(
292    post,
293    path = "/phases/{phase}/pause",
294    operation_id = "adminPausePhase",
295    tag = "credit-registration-admin",
296    params(("phase" = String, Path, description = "A canonical phase name")),
297    request_body = AdminPausePhasePayload,
298    responses(
299        (status = 200, description = "The phase's status after pausing", body = CreditRegistrationPhaseStatus),
300        (status = 422, description = "No reason given, or not one of the canonical phase names")
301    )
302)]
303pub async fn admin_pause_phase(
304    user: AuthUser,
305    pool: web::Data<PgPool>,
306    phase: web::Path<String>,
307    payload: web::Json<AdminPausePhasePayload>,
308) -> ControllerResult<web::Json<CreditRegistrationPhaseStatus>> {
309    let mut conn = pool.acquire().await?;
310    let token = authorize_credit_registration_admin(&mut conn, user.id).await?;
311
312    let phase = require_known_phase(&phase)?;
313    let reason = required_reason(&payload.reason)?;
314
315    info!(phase, actor = %user.id, "Admin paused credit registration phase");
316    let mut tx = conn.begin().await?;
317    credit_registration_phase_state::pause(&mut tx, phase, user.id, Some(reason)).await?;
318    record_phase_action(
319        &mut tx,
320        phase,
321        CreditRegistrationAdminAction::PausePhase,
322        user.id,
323        Some(reason.to_string()),
324    )
325    .await?;
326    tx.commit().await?;
327
328    token.authorized_ok(web::Json(one_phase_status(&mut conn, phase).await?))
329}
330
331/**
332POST `/api/v0/main-frontend/credit-registration-admin/phases/{phase}/resume` - Resumes one paused
333phase.
334*/
335#[instrument(skip(pool, payload))]
336#[utoipa::path(
337    post,
338    path = "/phases/{phase}/resume",
339    operation_id = "adminResumePhase",
340    tag = "credit-registration-admin",
341    params(("phase" = String, Path, description = "A canonical phase name")),
342    request_body = AdminPhaseActionPayload,
343    responses(
344        (status = 200, description = "The phase's status after resuming", body = CreditRegistrationPhaseStatus),
345        (status = 422, description = "Not one of the canonical phase names")
346    )
347)]
348pub async fn admin_resume_phase(
349    user: AuthUser,
350    pool: web::Data<PgPool>,
351    phase: web::Path<String>,
352    payload: web::Json<AdminPhaseActionPayload>,
353) -> ControllerResult<web::Json<CreditRegistrationPhaseStatus>> {
354    let mut conn = pool.acquire().await?;
355    let token = authorize_credit_registration_admin(&mut conn, user.id).await?;
356
357    let phase = require_known_phase(&phase)?;
358
359    info!(phase, actor = %user.id, "Admin resumed credit registration phase");
360    let mut tx = conn.begin().await?;
361    credit_registration_phase_state::resume(&mut tx, phase).await?;
362    record_phase_action(
363        &mut tx,
364        phase,
365        CreditRegistrationAdminAction::ResumePhase,
366        user.id,
367        payload.reason.clone(),
368    )
369    .await?;
370    tx.commit().await?;
371
372    token.authorized_ok(web::Json(one_phase_status(&mut conn, phase).await?))
373}
374
375/**
376POST `/api/v0/main-frontend/credit-registration-admin/phases/{phase}/run-now` - Makes one phase due
377immediately: the worker loop picks it up on its next tick instead of waiting out `next_run_at`.
378*/
379#[instrument(skip(pool, payload))]
380#[utoipa::path(
381    post,
382    path = "/phases/{phase}/run-now",
383    operation_id = "adminRunPhaseNow",
384    tag = "credit-registration-admin",
385    params(("phase" = String, Path, description = "A canonical phase name")),
386    request_body = AdminPhaseActionPayload,
387    responses(
388        (status = 200, description = "The phase's status after being made due", body = CreditRegistrationPhaseStatus),
389        (status = 422, description = "Not one of the canonical phase names")
390    )
391)]
392pub async fn admin_run_phase_now(
393    user: AuthUser,
394    pool: web::Data<PgPool>,
395    phase: web::Path<String>,
396    payload: web::Json<AdminPhaseActionPayload>,
397) -> ControllerResult<web::Json<CreditRegistrationPhaseStatus>> {
398    let mut conn = pool.acquire().await?;
399    let token = authorize_credit_registration_admin(&mut conn, user.id).await?;
400
401    let phase = require_known_phase(&phase)?;
402
403    info!(phase, actor = %user.id, "Admin forced credit registration phase to run now");
404    let mut tx = conn.begin().await?;
405    credit_registration_phase_state::run_now(&mut tx, phase).await?;
406    record_phase_action(
407        &mut tx,
408        phase,
409        CreditRegistrationAdminAction::RunPhaseNow,
410        user.id,
411        payload.reason.clone(),
412    )
413    .await?;
414    tx.commit().await?;
415
416    token.authorized_ok(web::Json(one_phase_status(&mut conn, phase).await?))
417}
418
419/// Records one admin action on a phase in the caller's transaction.
420async fn record_phase_action(
421    tx: &mut PgConnection,
422    phase: &str,
423    action: CreditRegistrationAdminAction,
424    actor_user_id: Uuid,
425    reason: Option<String>,
426) -> Result<(), ControllerError> {
427    models::credit_registration_admin_actions::record(
428        tx,
429        &NewCreditRegistrationAdminAction {
430            target_phase: Some(phase.to_string()),
431            reason,
432            ..NewCreditRegistrationAdminAction::new(
433                action,
434                CreditRegistrationAdminActionTarget::Phase,
435                actor_user_id,
436                GLOBAL_ADMIN_ROLE,
437            )
438        },
439    )
440    .await?;
441    Ok(())
442}
443
444/// Resolves a path segment to the spelling `credit_registration_phase_state` stores, refusing anything
445/// that is not a canonical phase name.
446fn require_known_phase(phase: &str) -> Result<&'static str, ControllerError> {
447    CreditRegistrationPhase::from_phase_name(phase)
448        .map(CreditRegistrationPhase::as_str)
449        .ok_or_else(|| {
450            controller_err!(
451                BadRequest,
452                "Not one of the canonical phase names.".to_string()
453            )
454        })
455}
456
457/// One phase's status, so a pause/resume/run-now response shows the effect without a second request.
458async fn one_phase_status(
459    conn: &mut PgConnection,
460    phase: &str,
461) -> Result<CreditRegistrationPhaseStatus, ControllerError> {
462    let row = credit_registration_phase_state::get_by_phase(conn, phase).await?;
463    Ok(to_phase_status(row, Utc::now()))
464}
465
466fn to_phase_status(
467    row: credit_registration_phase_state::CreditRegistrationPhaseState,
468    now: DateTime<Utc>,
469) -> CreditRegistrationPhaseStatus {
470    let seconds_since_heartbeat = row.last_heartbeat_at.map(|at| (now - at).num_seconds());
471    let heartbeat_late = is_heartbeat_late(
472        row.last_heartbeat_at,
473        row.expected_interval_secs,
474        row.paused_at,
475        now,
476    );
477    CreditRegistrationPhaseStatus {
478        is_known_phase: CreditRegistrationPhase::from_phase_name(&row.phase).is_some(),
479        phase: row.phase,
480        process_name: row.process_name,
481        expected_interval_secs: row.expected_interval_secs,
482        last_heartbeat_at: row.last_heartbeat_at,
483        last_success_at: row.last_success_at,
484        last_run_finished_at: row.last_run_finished_at,
485        items_processed_last_run: row.items_processed_last_run,
486        items_failed_last_run: row.items_failed_last_run,
487        consecutive_failures: row.consecutive_failures,
488        paused_at: row.paused_at,
489        pause_reason: row.pause_reason,
490        seconds_since_heartbeat,
491        heartbeat_late,
492    }
493}
494
495fn to_error_code_total(row: CreditRegistrationErrorCodeCount) -> CreditRegistrationErrorCodeTotal {
496    CreditRegistrationErrorCodeTotal {
497        error_code: row.error_code,
498        in_flight_count: row.in_flight_count,
499        terminal_failure_count: row.terminal_failure_count,
500    }
501}
502
503fn to_oldest_non_terminal(
504    row: OldestNonTerminalRegistration,
505    now: DateTime<Utc>,
506) -> CreditRegistrationOldestNonTerminal {
507    CreditRegistrationOldestNonTerminal {
508        credit_registration_id: row.id,
509        state: row.state,
510        seconds_in_state: (now - row.state_entered_at).num_seconds(),
511        state_entered_at: row.state_entered_at,
512    }
513}
514
515fn to_stuck_total(row: StuckRegistrationCount) -> CreditRegistrationStuckTotal {
516    CreditRegistrationStuckTotal {
517        state: row.state,
518        count: row.count,
519        severely_stuck_count: row.severely_stuck_count,
520        oldest_state_entered_at: row.oldest_state_entered_at,
521    }
522}
523
524fn to_endpoint_standing(row: SuotarEndpointStandingRow) -> SuotarEndpointStanding {
525    SuotarEndpointStanding {
526        endpoint: row.endpoint,
527        last_success_at: row.last_success_at,
528        last_failure_at: row.last_failure_at,
529        consecutive_failures: row.consecutive_failures,
530    }
531}
532
533fn to_endpoint_window_stats_for_window(
534    row: SuotarEndpointStatsForWindow,
535) -> SuotarEndpointWindowStats {
536    SuotarEndpointWindowStats {
537        endpoint: row.endpoint,
538        call_count: row.call_count,
539        failed_call_count: row.failed_call_count,
540        in_flight_count: row.in_flight_count,
541        ok_item_count: row.ok_item_count,
542        error_item_count: row.error_item_count,
543        pending_item_count: row.pending_item_count,
544        p50_duration_ms: row.p50_duration_ms,
545        p95_duration_ms: row.p95_duration_ms,
546        last_success_at: row.last_success_at,
547        last_failure_at: row.last_failure_at,
548        last_request_level_error_code: row.last_request_level_error_code,
549    }
550}
551
552pub fn _add_routes(cfg: &mut ServiceConfig) {
553    cfg.route("/overview", web::get().to(get_credit_registration_overview))
554        .route("/suotar-health", web::get().to(get_suotar_health))
555        .route("/phases/{phase}/pause", web::post().to(admin_pause_phase))
556        .route("/phases/{phase}/resume", web::post().to(admin_resume_phase))
557        .route(
558            "/phases/{phase}/run-now",
559            web::post().to(admin_run_phase_now),
560        );
561}