Skip to main content

headless_lms_credit_registration/runtime/
worker_loop.rs

1//! The loop both credit registration workers run.
2//!
3//! The processes differ only in which phases they own and how often they look, so the scheduling
4//! lives here instead of in each of them.
5
6use std::{
7    sync::Arc,
8    time::{Duration, Instant},
9};
10
11use chrono::{DateTime, Utc};
12use sqlx::PgPool;
13
14use headless_lms_base::config::ApplicationConfiguration;
15use headless_lms_models::credit_registration_phase_state::{
16    self, CreditRegistrationPhaseState, set_next_run_at,
17};
18use headless_lms_models::suotar_api_calls::PgSuotarCallAudit;
19use headless_lms_utils::services::suotar::SuotarClient;
20use tokio_util::sync::CancellationToken;
21
22use super::is_waiting_item;
23
24use headless_lms_utils::periodic_worker::{
25    PeriodicWorkerConfig, StillRunningLog, run_periodic_worker_until,
26};
27
28use super::dispatch::{PhaseContext, PhaseSkipReason, PhaseTick, Runner, run_phase_once};
29use super::process_local::LastReported;
30use crate::error::{CreditRegistrationError, CreditRegistrationResult};
31use crate::error_reports::ErrorReporter;
32use crate::phase::{CreditRegistrationPhase, WorkerProcess};
33use headless_lms_models::credit_registrations::RegistrationScope;
34
35/// How often each phase's loop looks whether it is due; each phase's own interval lives in
36/// `credit_registration_phase_state`.
37const TICK_INTERVAL: Duration = Duration::from_secs(10);
38
39/// Ten minutes of ticks. The per-phase heartbeat in the database is the machine-readable half.
40const STILL_RUNNING_MESSAGE_TICKS: u32 = 60;
41
42/// Runs the phases `process` owns until SIGTERM or Ctrl-C. Each phase loops on its own, so an
43/// hour-long call in one does not hold up the others. On shutdown no phase starts another
44/// iteration, and the function returns once the iterations already running have finished.
45pub async fn run(
46    process: WorkerProcess,
47    db_pool: PgPool,
48    app_configuration: ApplicationConfiguration,
49    still_running_message: &str,
50) -> CreditRegistrationResult<()> {
51    let suotar_client = SuotarClient::new(
52        &app_configuration.suotar_configuration,
53        Arc::new(PgSuotarCallAudit::new(db_pool.clone(), is_waiting_item)),
54    );
55    let shutdown = CancellationToken::new();
56    tokio::spawn(cancel_on_termination_signal(shutdown.clone()));
57    let ctx = PhaseContext {
58        shutdown: Some(&shutdown),
59        ..PhaseContext::from_app(
60            &db_pool,
61            &suotar_client,
62            &app_configuration,
63            Runner::Worker(process),
64        )
65    };
66
67    // Its body is empty on purpose: the helper's own still-running line is all this loop is for.
68    let still_running = run_periodic_worker_until(
69        PeriodicWorkerConfig {
70            tick_interval: TICK_INTERVAL,
71            still_running: Some(StillRunningLog {
72                every: STILL_RUNNING_MESSAGE_TICKS,
73                message: still_running_message,
74                initial_ticks: 0,
75            }),
76            delay_missed_ticks: true,
77        },
78        &shutdown,
79        async || Ok::<(), CreditRegistrationError>(()),
80    );
81    // Futures of one task rather than spawned tasks: the phase bodies are not `Send`. They only
82    // need to wait on the study registry side by side, not to run in parallel.
83    let phase_loops = CreditRegistrationPhase::ALL
84        .into_iter()
85        .filter(|phase| phase.spec().process == process)
86        .map(|phase| run_phase_loop(&ctx, phase, &shutdown));
87    let (still_running, phase_loops) =
88        tokio::join!(still_running, futures::future::join_all(phase_loops));
89    still_running?;
90    phase_loops
91        .into_iter()
92        .collect::<CreditRegistrationResult<()>>()?;
93    info!("{} stopped.", process.as_str());
94    Ok(())
95}
96
97async fn run_phase_loop(
98    ctx: &PhaseContext<'_>,
99    phase: CreditRegistrationPhase,
100    shutdown: &CancellationToken,
101) -> CreditRegistrationResult<()> {
102    run_periodic_worker_until(
103        PeriodicWorkerConfig {
104            tick_interval: TICK_INTERVAL,
105            // `run` logs one message for the whole process.
106            still_running: None,
107            // A slow iteration should push later ticks out, not fire them back to back (tokio's
108            // default).
109            delay_missed_ticks: true,
110        },
111        shutdown,
112        async || {
113            trace!(phase = phase.as_str(), "Checking whether phase is due");
114            // Logged and swallowed: the phase-state row already carries the failure for the
115            // dashboard, and the loop must keep going.
116            if let Err(error) = run_if_due(ctx, phase).await {
117                log_failure(ctx.runner.caller(), phase.as_str(), &error);
118                ErrorReporter::new(ctx.pool, ctx.runner.owning_process(), phase)
119                    .report(
120                        &error.cause_chain(),
121                        Some(format!("{error:?}")),
122                        serde_json::json!({}),
123                    )
124                    .await;
125            }
126            Ok(())
127        },
128    )
129    .await
130}
131
132/// Kubernetes sends SIGTERM and waits `terminationGracePeriodSeconds` before killing the pod.
133async fn cancel_on_termination_signal(shutdown: CancellationToken) {
134    let terminate = async {
135        match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
136            Ok(mut signal) => {
137                signal.recv().await;
138            }
139            Err(error) => {
140                error!(error = %error, "Could not listen for SIGTERM");
141                std::future::pending::<()>().await;
142            }
143        }
144    };
145    tokio::select! {
146        _ = terminate => info!("Received SIGTERM; finishing the phase iterations already running."),
147        _ = tokio::signal::ctrl_c() => info!("Received Ctrl-C; finishing the phase iterations already running."),
148    }
149    shutdown.cancel();
150}
151
152async fn run_if_due(
153    ctx: &PhaseContext<'_>,
154    phase: CreditRegistrationPhase,
155) -> CreditRegistrationResult<()> {
156    let state = {
157        let mut conn = ctx.pool.acquire().await?;
158        credit_registration_phase_state::get_by_phase(&mut conn, phase.as_str()).await?
159    };
160    if !is_due(&state, Utc::now()) {
161        return Ok(());
162    }
163    run_due_phase(ctx, phase, &state).await
164}
165
166fn log_failure(process_name: &str, phase: &str, error: &CreditRegistrationError) {
167    error!(phase, error = %error, "Credit registration phase iteration failed");
168    if error.is_db_disconnect() {
169        info!(process_name, "May have lost its connection to the database");
170    }
171}
172
173/// Runs one due phase, and schedules the next run.
174async fn run_due_phase(
175    ctx: &PhaseContext<'_>,
176    phase: CreditRegistrationPhase,
177    state: &CreditRegistrationPhaseState,
178) -> CreditRegistrationResult<()> {
179    let mut conn = ctx.pool.acquire().await?;
180    // Stamped before the work, so a phase whose iteration takes longer than its interval does not
181    // run back to back.
182    set_next_run_at(
183        &mut conn,
184        phase.as_str(),
185        Utc::now() + chrono::Duration::seconds(state.expected_interval_secs.into()),
186    )
187    .await?;
188    drop(conn);
189
190    let started_at = Instant::now();
191    // Always unscoped: a worker that narrowed would leave rows nobody sweeps.
192    let tick = run_phase_once(ctx, phase, &RegistrationScope::default()).await?;
193    let duration_ms = started_at.elapsed().as_millis() as u64;
194    match tick {
195        PhaseTick::Ran(outcome)
196            if outcome.items_processed > 0
197                || outcome.items_waiting > 0
198                || outcome.items_failed > 0 =>
199        {
200            clear_skip_state(phase);
201            let processed = outcome.items_processed;
202            let waiting = outcome.items_waiting;
203            let failed = outcome.items_failed;
204            info!(
205                phase = phase.as_str(),
206                processed,
207                waiting,
208                failed,
209                duration_ms,
210                "processed {processed} rows ({waiting} waiting, {failed} failed), took {duration_ms}ms"
211            );
212        }
213        PhaseTick::Ran(_) => {
214            clear_skip_state(phase);
215            // Nothing to do this run: phases tick every few seconds, so anything above trace buries
216            // the lines that report work.
217            trace!(
218                phase = phase.as_str(),
219                duration_ms, "Credit registration phase run found nothing to do"
220            );
221        }
222        // Waiting out a cooldown, or turned off: the phase-state row already carries this for the
223        // dashboard, so only the state change is worth a log line, not every retry.
224        PhaseTick::Skipped(reason) => log_skip_if_changed(phase, reason),
225        PhaseTick::ScopeNotSupported => {}
226    }
227    Ok(())
228}
229
230/// The skip reason last logged for a phase, so a paused phase or an open breaker logs once per
231/// state change instead of on every tick until it clears.
232static LAST_LOGGED_SKIP: LastReported<CreditRegistrationPhase, PhaseSkipReason> =
233    LastReported::without_refresh();
234
235fn log_skip_if_changed(phase: CreditRegistrationPhase, reason: PhaseSkipReason) {
236    if !LAST_LOGGED_SKIP.is_due(&phase, &reason, PartialEq::eq) {
237        return;
238    }
239    LAST_LOGGED_SKIP.record(phase, reason);
240    match reason {
241        PhaseSkipReason::Paused => {
242            info!(
243                phase = phase.as_str(),
244                "Credit registration phase is paused; skipping"
245            );
246        }
247        PhaseSkipReason::CircuitBreakerOpen => {
248            info!(
249                phase = phase.as_str(),
250                "Credit registration phase skipped: circuit breaker is open"
251            );
252        }
253        PhaseSkipReason::AccountLinkingDisabled => {
254            debug!(
255                phase = phase.as_str(),
256                "Credit registration phase skipped: account linking is disabled"
257            );
258        }
259        PhaseSkipReason::SisuDayGap => {
260            info!(
261                phase = phase.as_str(),
262                "Credit registration phase skipped: Sisu is still on the previous day"
263            );
264        }
265    }
266}
267
268fn clear_skip_state(phase: CreditRegistrationPhase) {
269    LAST_LOGGED_SKIP.clear(&phase);
270}
271
272/// A phase is due when an admin asked for it, or when its interval has elapsed since it last began.
273fn is_due(state: &CreditRegistrationPhaseState, now: DateTime<Utc>) -> bool {
274    if let Some(next_run_at) = state.next_run_at {
275        return next_run_at <= now;
276    }
277    state.last_run_started_at.is_none_or(|started| {
278        (now - started).num_seconds() >= i64::from(state.expected_interval_secs)
279    })
280}
281
282#[cfg(test)]
283mod tests {
284    use super::*;
285    use uuid::Uuid;
286
287    fn state(
288        next_run_at: Option<DateTime<Utc>>,
289        last_run_started_at: Option<DateTime<Utc>>,
290    ) -> CreditRegistrationPhaseState {
291        CreditRegistrationPhaseState {
292            id: Uuid::new_v4(),
293            created_at: Utc::now(),
294            updated_at: Utc::now(),
295            deleted_at: None,
296            phase: "verify".to_string(),
297            process_name: "credit-registrar".to_string(),
298            expected_interval_secs: 60,
299            last_heartbeat_at: None,
300            last_run_started_at,
301            last_run_finished_at: None,
302            last_success_at: None,
303            next_run_at,
304            items_processed_last_run: None,
305            items_failed_last_run: None,
306            consecutive_failures: 0,
307            last_error: None,
308            paused_at: None,
309            paused_by_user_id: None,
310            pause_reason: None,
311        }
312    }
313
314    #[test]
315    fn a_phase_that_has_never_run_is_due() {
316        assert!(is_due(&state(None, None), Utc::now()));
317    }
318
319    #[test]
320    fn a_phase_is_due_again_once_its_interval_has_elapsed() {
321        let now = Utc::now();
322        assert!(!is_due(
323            &state(None, Some(now - chrono::Duration::seconds(30))),
324            now
325        ));
326        assert!(is_due(
327            &state(None, Some(now - chrono::Duration::seconds(90))),
328            now
329        ));
330    }
331
332    /// How "run now" works: the admin endpoint stamps the timestamp and the loop notices.
333    #[test]
334    fn an_explicit_next_run_beats_the_interval() {
335        let now = Utc::now();
336        let asked_for = state(Some(now), Some(now));
337        assert!(is_due(&asked_for, now));
338
339        let scheduled = state(Some(now + chrono::Duration::seconds(30)), None);
340        assert!(!is_due(&scheduled, now));
341    }
342
343    /// A phase belonging to neither process would look merely idle rather than unrun.
344    #[test]
345    fn the_two_processes_between_them_own_every_phase() {
346        let mut owned: Vec<&str> = Vec::new();
347        for process in WorkerProcess::ALL {
348            owned.extend(
349                CreditRegistrationPhase::ALL
350                    .into_iter()
351                    .filter(|phase| phase.spec().process == process)
352                    .map(|phase| phase.as_str()),
353            );
354        }
355        assert_eq!(owned.len(), CreditRegistrationPhase::ALL.len());
356    }
357}