1use 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 pub in_flight_count: i64,
44 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 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 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#[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#[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 pub is_known_phase: bool,
103 pub seconds_since_heartbeat: Option<i64>,
106 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 pub pending_by_reason: PendingReasonCounts,
117 pub error_codes: Vec<CreditRegistrationErrorCodeTotal>,
118 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 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#[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#[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#[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#[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#[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
419async 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
444fn 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
457async 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}