Skip to main content

headless_lms_models/
credit_registration_phase_state.rs

1//! Per-phase heartbeat and control for the credit registration pipeline.
2//!
3//! One row per phase, seeded by migration and thereafter only ever updated.
4use utoipa::ToSchema;
5
6use crate::prelude::*;
7
8/// Canonical phase names, used verbatim as the `phase` value, in the test tick endpoint, the
9/// dashboard and the audit log.
10pub const PHASES: &[&str] = &[
11    "materialize",
12    "preconditions",
13    "resolve-enrolments",
14    "import",
15    "verify",
16    "legacy-mirror",
17    "student-notifications",
18    "enrolment-discovery",
19    "link-emails",
20    "config-validation",
21    "retention-sweep",
22    "ledger-snapshot",
23];
24
25#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, ToSchema)]
26pub struct CreditRegistrationPhaseState {
27    pub id: Uuid,
28    pub created_at: DateTime<Utc>,
29    pub updated_at: DateTime<Utc>,
30    pub deleted_at: Option<DateTime<Utc>>,
31    pub phase: String,
32    pub process_name: String,
33    pub expected_interval_secs: i32,
34    pub last_heartbeat_at: Option<DateTime<Utc>>,
35    pub last_run_started_at: Option<DateTime<Utc>>,
36    pub last_run_finished_at: Option<DateTime<Utc>>,
37    pub last_success_at: Option<DateTime<Utc>>,
38    pub next_run_at: Option<DateTime<Utc>>,
39    pub items_processed_last_run: Option<i32>,
40    pub items_failed_last_run: Option<i32>,
41    pub consecutive_failures: i32,
42    pub last_error: Option<String>,
43    pub paused_at: Option<DateTime<Utc>>,
44    pub paused_by_user_id: Option<Uuid>,
45    pub pause_reason: Option<String>,
46}
47
48/// What one iteration of a phase did, as its phase-state row records it.
49#[derive(Debug, Clone, PartialEq, Default)]
50pub struct PhaseRunOutcome {
51    pub items_processed: i32,
52    /// Processed rows that only await something outside the pipeline; not in `items_failed`, and
53    /// not stored.
54    pub items_waiting: i32,
55    pub items_failed: i32,
56    /// `None` on success. Scrub before passing.
57    pub error: Option<String>,
58}
59
60pub async fn get_all(conn: &mut PgConnection) -> ModelResult<Vec<CreditRegistrationPhaseState>> {
61    let res = sqlx::query_as!(
62        CreditRegistrationPhaseState,
63        r#"
64SELECT *
65FROM credit_registration_phase_state
66WHERE deleted_at IS NULL
67ORDER BY process_name,
68  phase
69        "#,
70    )
71    .fetch_all(conn)
72    .await?;
73    Ok(res)
74}
75
76pub async fn get_by_phase(
77    conn: &mut PgConnection,
78    phase: &str,
79) -> ModelResult<CreditRegistrationPhaseState> {
80    let res = sqlx::query_as!(
81        CreditRegistrationPhaseState,
82        r#"
83SELECT *
84FROM credit_registration_phase_state
85WHERE phase = $1
86  AND deleted_at IS NULL
87        "#,
88        phase
89    )
90    .fetch_one(conn)
91    .await?;
92    Ok(res)
93}
94
95/// Written every iteration, work or not, so idle and wedged stay distinguishable.
96pub async fn heartbeat(conn: &mut PgConnection, phase: &str) -> ModelResult<()> {
97    sqlx::query!(
98        r#"
99UPDATE credit_registration_phase_state
100SET last_heartbeat_at = now(),
101  last_run_started_at = now()
102WHERE phase = $1
103  AND deleted_at IS NULL
104        "#,
105        phase
106    )
107    .execute(conn)
108    .await?;
109    Ok(())
110}
111
112/// Keeps a long iteration from reading as a dead worker; unlike [`heartbeat`], leaves
113/// `last_run_started_at` at the iteration's start.
114pub async fn keep_alive(conn: &mut PgConnection, phase: &str) -> ModelResult<()> {
115    sqlx::query!(
116        r#"
117UPDATE credit_registration_phase_state
118SET last_heartbeat_at = now()
119WHERE phase = $1
120  AND deleted_at IS NULL
121        "#,
122        phase
123    )
124    .execute(conn)
125    .await?;
126    Ok(())
127}
128
129/// Closes out an iteration. Only a success moves `last_success_at`, so a wedged phase stays
130/// distinguishable from a quiet one.
131pub async fn record_run(
132    conn: &mut PgConnection,
133    phase: &str,
134    outcome: &PhaseRunOutcome,
135) -> ModelResult<()> {
136    sqlx::query!(
137        r#"
138UPDATE credit_registration_phase_state
139SET last_run_finished_at = now(),
140  items_processed_last_run = $2,
141  items_failed_last_run = $3,
142  last_success_at = CASE
143    WHEN $4::text IS NULL THEN now()
144    ELSE last_success_at
145  END,
146  consecutive_failures = CASE
147    WHEN $4::text IS NULL THEN 0
148    ELSE consecutive_failures + 1
149  END,
150  last_error = $4
151WHERE phase = $1
152  AND deleted_at IS NULL
153        "#,
154        phase,
155        outcome.items_processed,
156        outcome.items_failed,
157        outcome.error,
158    )
159    .execute(conn)
160    .await?;
161    Ok(())
162}
163
164pub async fn is_paused(conn: &mut PgConnection, phase: &str) -> ModelResult<bool> {
165    let paused = sqlx::query_scalar!(
166        r#"
167SELECT paused_at IS NOT NULL AS "paused!"
168FROM credit_registration_phase_state
169WHERE phase = $1
170  AND deleted_at IS NULL
171        "#,
172        phase
173    )
174    .fetch_one(conn)
175    .await?;
176    Ok(paused)
177}
178
179pub async fn pause(
180    conn: &mut PgConnection,
181    phase: &str,
182    paused_by_user_id: Uuid,
183    pause_reason: Option<&str>,
184) -> ModelResult<()> {
185    sqlx::query!(
186        r#"
187UPDATE credit_registration_phase_state
188SET paused_at = now(),
189  paused_by_user_id = $2,
190  pause_reason = $3
191WHERE phase = $1
192  AND deleted_at IS NULL
193        "#,
194        phase,
195        paused_by_user_id,
196        pause_reason,
197    )
198    .execute(conn)
199    .await?;
200    Ok(())
201}
202
203pub async fn resume(conn: &mut PgConnection, phase: &str) -> ModelResult<()> {
204    sqlx::query!(
205        r#"
206UPDATE credit_registration_phase_state
207SET paused_at = NULL,
208  paused_by_user_id = NULL,
209  pause_reason = NULL
210WHERE phase = $1
211  AND deleted_at IS NULL
212        "#,
213        phase
214    )
215    .execute(conn)
216    .await?;
217    Ok(())
218}
219
220/// Makes the phase due now; the phase loop picks it up on its next `next_run_at` check.
221pub async fn run_now(conn: &mut PgConnection, phase: &str) -> ModelResult<()> {
222    sqlx::query!(
223        r#"
224UPDATE credit_registration_phase_state
225SET next_run_at = now()
226WHERE phase = $1
227  AND deleted_at IS NULL
228        "#,
229        phase
230    )
231    .execute(conn)
232    .await?;
233    Ok(())
234}
235
236pub async fn set_next_run_at(
237    conn: &mut PgConnection,
238    phase: &str,
239    next_run_at: DateTime<Utc>,
240) -> ModelResult<()> {
241    sqlx::query!(
242        r#"
243UPDATE credit_registration_phase_state
244SET next_run_at = $2
245WHERE phase = $1
246  AND deleted_at IS NULL
247        "#,
248        phase,
249        next_run_at,
250    )
251    .execute(conn)
252    .await?;
253    Ok(())
254}