Skip to main content

headless_lms_models/
suotar_circuit_breakers.rs

1//! The last state each worker process reported for its circuit breakers. The worker keeps the live
2//! state in memory; this copy is only for the dashboard, which runs in another process.
3
4use utoipa::ToSchema;
5
6use crate::prelude::*;
7
8/// Which phases one circuit breaker pauses.
9#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, Type, ToSchema)]
10#[sqlx(type_name = "suotar_circuit_breaker_target", rename_all = "snake_case")]
11#[serde(rename_all = "snake_case")]
12pub enum BreakerTarget {
13    /// Every phase that calls the study registry: Suotar itself failing.
14    StudyRegistry,
15    /// Only the phase that submits to Sisu: Suotar answering that Sisu timed out.
16    SisuSubmissions,
17}
18
19/// One worker process's circuit breaker as that process last reported it, for the dashboard. Stale
20/// once the process stops reporting: check `updated_at`.
21#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, ToSchema)]
22pub struct SuotarCircuitBreaker {
23    /// The worker process the breaker belongs to; each keeps its own.
24    pub process_name: String,
25    pub target: BreakerTarget,
26    /// When the worker last reported the state.
27    pub updated_at: DateTime<Utc>,
28    pub consecutive_failures: i32,
29    /// When the cooldown ends; `None` for a breaker that has not opened since its last success.
30    pub open_until: Option<DateTime<Utc>>,
31    pub trip_count: i32,
32}
33
34/// What the worker reports for one breaker; [`SuotarCircuitBreaker`] without the report time.
35#[derive(Debug, Clone, PartialEq)]
36pub struct SuotarCircuitBreakerReport<'a> {
37    pub process_name: &'a str,
38    pub target: BreakerTarget,
39    pub consecutive_failures: i32,
40    pub open_until: Option<DateTime<Utc>>,
41    pub trip_count: i32,
42}
43
44/// Replaces the reported state of one process's breaker. Only the worker process that owns the
45/// breaker may call this: any other process holds a breaker of its own, and would overwrite the
46/// worker's.
47pub async fn upsert(
48    conn: &mut PgConnection,
49    breaker: &SuotarCircuitBreakerReport<'_>,
50) -> ModelResult<()> {
51    sqlx::query!(
52        r#"
53INSERT INTO suotar_circuit_breakers (
54    process_name,
55    target,
56    consecutive_failures,
57    open_until,
58    trip_count
59  )
60VALUES ($1, $2, $3, $4, $5) ON CONFLICT (process_name, target) DO
61UPDATE
62SET consecutive_failures = EXCLUDED.consecutive_failures,
63  open_until = EXCLUDED.open_until,
64  trip_count = EXCLUDED.trip_count
65        "#,
66        breaker.process_name,
67        breaker.target as BreakerTarget,
68        breaker.consecutive_failures,
69        breaker.open_until,
70        breaker.trip_count,
71    )
72    .execute(conn)
73    .await?;
74    Ok(())
75}
76
77/// Every reported breaker, by process and target.
78pub async fn get_all(conn: &mut PgConnection) -> ModelResult<Vec<SuotarCircuitBreaker>> {
79    let res = sqlx::query_as!(
80        SuotarCircuitBreaker,
81        r#"
82SELECT process_name,
83  target,
84  updated_at,
85  consecutive_failures,
86  open_until,
87  trip_count
88FROM suotar_circuit_breakers
89ORDER BY process_name,
90  target
91        "#,
92    )
93    .fetch_all(conn)
94    .await?;
95    Ok(res)
96}