Skip to main content

headless_lms_credit_registration/runtime/
dispatch.rs

1//! Running one iteration of one phase, and the bookkeeping around it: the pause and scope checks,
2//! the heartbeat, the circuit breakers and the limiter.
3
4use chrono::TimeDelta;
5use headless_lms_models::credit_registration_phase_state::{self, PhaseRunOutcome};
6use headless_lms_models::library::credit_registration::scrub::scrub_text;
7use headless_lms_models::library::credit_registration::sisu_day_gap;
8use headless_lms_utils::services::suotar::SuotarClient;
9use sqlx::PgPool;
10use tokio_util::sync::CancellationToken;
11
12use super::heartbeat::keep_alive;
13use super::suotar::{
14    SuotarStudyRegistry, max_study_registry_wait, report_breakers, report_rate_limits,
15};
16use crate::error::CreditRegistrationResult;
17use crate::error_reports::ErrorReporter;
18use crate::phase::{CreditRegistrationPhase, WorkerProcess};
19use crate::use_cases::batch_flow::BatchFlowContext;
20use crate::use_cases::{
21    config_validation, enrolment_discovery, import, ledger_snapshot, legacy_mirror, link_emails,
22    materialize, preconditions, resolve_enrolments, retention_sweep, student_notifications, verify,
23};
24use crate::workflow::Counts;
25use headless_lms_models::credit_registrations::RegistrationScope;
26
27/// What one dispatch attempt did.
28#[derive(Debug, Clone, PartialEq)]
29pub enum PhaseTick {
30    Ran(PhaseRunOutcome),
31    /// The phase legitimately did nothing; not counted as a failure.
32    Skipped(PhaseSkipReason),
33    /// The scope names something this phase cannot narrow on; refused rather than run wide.
34    ScopeNotSupported,
35}
36
37#[derive(Debug, Clone, Copy, PartialEq, Eq)]
38pub enum PhaseSkipReason {
39    Paused,
40    CircuitBreakerOpen,
41    AccountLinkingDisabled,
42    /// Finland is on the next day and Sisu is not; see [`sisu_day_gap`].
43    SisuDayGap,
44}
45
46/// Who runs a phase iteration.
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub enum Runner<'a> {
49    /// A worker process's loop, the only holder of the breaker and limiter state the dashboard
50    /// shows for the phases it owns.
51    Worker(WorkerProcess),
52    /// Anyone else, such as the test tick endpoint, named for the audit log. Its in-memory state
53    /// would overwrite the worker's on the dashboard, so it reports none.
54    Other(&'a str),
55}
56
57impl Runner<'_> {
58    /// What goes into the audit log's `worker_name` alongside the phase.
59    pub fn caller(&self) -> &str {
60        match self {
61            Self::Worker(process) => process.as_str(),
62            Self::Other(caller) => caller,
63        }
64    }
65
66    pub fn owning_process(&self) -> Option<WorkerProcess> {
67        match self {
68            Self::Worker(process) => Some(*process),
69            Self::Other(_) => None,
70        }
71    }
72}
73
74/// Everything a phase iteration needs from its caller: the worker loop or the test tick endpoint.
75pub struct PhaseContext<'a> {
76    pub pool: &'a PgPool,
77    pub suotar_client: &'a SuotarClient,
78    /// Shortens the circuit breaker's cooldown to something a test can wait out.
79    pub test_mode: bool,
80    pub runner: Runner<'a>,
81    /// Absolute base for links in queued mail, which outlive the process that wrote them.
82    pub base_url: &'a str,
83    /// Off, the linking mails are not sent, and discovery only wakes linked students' registrations.
84    pub is_account_linking_enabled: bool,
85    /// The worker's SIGTERM; `None` for a run no signal can stop, such as an on-demand one.
86    pub shutdown: Option<&'a CancellationToken>,
87}
88
89impl<'a> PhaseContext<'a> {
90    /// Builds a context from the application configuration, the shape every construction site
91    /// starts from.
92    pub fn from_app(
93        pool: &'a PgPool,
94        suotar_client: &'a SuotarClient,
95        app_conf: &'a headless_lms_base::config::ApplicationConfiguration,
96        runner: Runner<'a>,
97    ) -> Self {
98        Self {
99            pool,
100            suotar_client,
101            test_mode: app_conf.test_mode,
102            runner,
103            base_url: &app_conf.base_url,
104            is_account_linking_enabled: app_conf.suotar_configuration.account_linking_enabled,
105            shutdown: None,
106        }
107    }
108}
109
110/// The audit log's `worker_name`, which the database caps at 64 characters.
111pub(super) fn worker_name(caller: &str, phase: CreditRegistrationPhase) -> String {
112    format!("{caller}/{}", phase.as_str())
113}
114
115/// Runs exactly one iteration of one phase.
116#[tracing::instrument(
117    skip_all,
118    fields(phase = phase.as_str(), caller = ctx.runner.caller(), scope = ?scope)
119)]
120pub async fn run_phase_once(
121    ctx: &PhaseContext<'_>,
122    phase: CreditRegistrationPhase,
123    scope: &RegistrationScope,
124) -> CreditRegistrationResult<PhaseTick> {
125    // Before the pause check: a caller whose narrowing cannot be honoured must not be told it ran.
126    if !phase.spec().scope.covers(scope) {
127        return Ok(PhaseTick::ScopeNotSupported);
128    }
129    let mut conn = ctx.pool.acquire().await?;
130    if credit_registration_phase_state::is_paused(&mut conn, phase.as_str()).await? {
131        return Ok(PhaseTick::Skipped(PhaseSkipReason::Paused));
132    }
133    // A scoped run writes nothing to the phase-state row: that row describes the workers, and a
134    // test's traffic in it would make a dead worker look alive to the heartbeat alert.
135    let bookkeeping = scope.is_unscoped();
136    let is_own_worker = ctx.runner.owning_process() == Some(phase.spec().process);
137    // Before the breaker check, unlike the pause above, which health.rs excludes from the staleness
138    // alert by itself. A cooldown is a worker deliberately waiting, not a worker that died, and
139    // skipping the heartbeat through it would raise a critical alert within a tick or two.
140    if bookkeeping {
141        credit_registration_phase_state::heartbeat(&mut conn, phase.as_str()).await?;
142        trace!(phase = phase.as_str(), "Wrote phase heartbeat");
143    }
144    // After the heartbeat, like the breaker check below: a switched-off phase is idle, not dead.
145    if phase.spec().is_account_linking_only && !ctx.is_account_linking_enabled {
146        return Ok(PhaseTick::Skipped(PhaseSkipReason::AccountLinkingDisabled));
147    }
148    // Not in test mode: the system tests run at every hour, the gap's included.
149    if phase.spec().waits_out_sisu_day_gap
150        && !ctx.test_mode
151        && let Some(gap_end) = sisu_day_gap::current_gap_end(&mut conn).await?
152    {
153        let moved = sisu_day_gap::spread_imports_past_gap(&mut conn, gap_end).await?;
154        if moved > 0 {
155            info!(
156                phase = phase.as_str(),
157                moved,
158                %gap_end,
159                "Spread imports due in the Sisu day gap over the hours after it"
160            );
161        }
162        if bookkeeping {
163            credit_registration_phase_state::record_deliberate_idle_run(&mut conn, phase.as_str())
164                .await?;
165            if is_own_worker {
166                report_breakers(&mut conn, phase).await?;
167            }
168        }
169        return Ok(PhaseTick::Skipped(PhaseSkipReason::SisuDayGap));
170    }
171    let mut registry = match SuotarStudyRegistry::admit(
172        ctx.suotar_client,
173        worker_name(ctx.runner.caller(), phase),
174        phase,
175        scope,
176        ctx.test_mode,
177    ) {
178        Ok(registry) => registry,
179        Err(skip) => {
180            if bookkeeping && is_own_worker {
181                report_breakers(&mut conn, phase).await?;
182            }
183            return Ok(PhaseTick::Skipped(skip));
184        }
185    };
186    drop(conn);
187
188    let body = if bookkeeping {
189        tokio::select! {
190            body = run_body(ctx, phase, scope, &mut registry) => body,
191            never = keep_alive(ctx.pool, phase) => match never {},
192        }
193    } else {
194        run_body(ctx, phase, scope, &mut registry).await
195    };
196    let failure = registry.finish();
197    // Only a real `CreditRegistrationError` carries a backtrace and span trace; `failure` and
198    // `counts.finding` are already plain messages.
199    let stack_trace = match &body {
200        Err(error) => Some(format!("{error:?}")),
201        Ok(_) => None,
202    };
203    let outcome = match body {
204        Ok(counts) => {
205            let outcome = PhaseRunOutcome {
206                items_processed: counts.processed_count(),
207                items_waiting: counts.waiting_count(),
208                items_failed: counts.failed_count(),
209                error: failure.or(counts.into_finding()),
210            };
211            if let Some(error) = &outcome.error {
212                warn!(phase = phase.as_str(), error = %error, "Credit registration phase iteration recorded a failure");
213            }
214            outcome
215        }
216        Err(error) => {
217            error!(phase = phase.as_str(), error = %error, "Credit registration phase iteration aborted");
218            PhaseRunOutcome {
219                error: Some(scrub_text(&error.cause_chain())),
220                ..PhaseRunOutcome::default()
221            }
222        }
223    };
224    if let Some(error) = &outcome.error {
225        ErrorReporter::new(ctx.pool, ctx.runner.owning_process(), phase)
226            .report(error, stack_trace, serde_json::json!({}))
227            .await;
228    }
229    if bookkeeping {
230        let mut conn = ctx.pool.acquire().await?;
231        credit_registration_phase_state::record_run(&mut conn, phase.as_str(), &outcome).await?;
232        if is_own_worker {
233            report_rate_limits(&mut conn, phase).await?;
234            report_breakers(&mut conn, phase).await?;
235        }
236    }
237    Ok(PhaseTick::Ran(outcome))
238}
239
240/// The one place a phase implementation is registered: exhaustive over [`CreditRegistrationPhase`],
241/// so a variant added there without an arm here fails to compile. Each phase gets only what it
242/// touches, and a phase that asks the study registry gets the iteration's registry beside it.
243async fn run_body(
244    ctx: &PhaseContext<'_>,
245    phase: CreditRegistrationPhase,
246    scope: &RegistrationScope,
247    registry: &mut SuotarStudyRegistry<'_>,
248) -> CreditRegistrationResult<Counts> {
249    let pool = ctx.pool;
250    let batch_flow = BatchFlowContext {
251        pool,
252        scope,
253        phase,
254        errors: ErrorReporter::new(pool, ctx.runner.owning_process(), phase),
255        shutdown: ctx.shutdown,
256        // Request timeouts are minutes, far inside the range.
257        study_registry_wait: TimeDelta::from_std(max_study_registry_wait(phase))
258            .unwrap_or(TimeDelta::zero()),
259    };
260    match phase {
261        CreditRegistrationPhase::Materialize => materialize::run(pool, scope).await,
262        CreditRegistrationPhase::Preconditions => preconditions::run(pool, scope).await,
263        CreditRegistrationPhase::ResolveEnrolments => {
264            resolve_enrolments::run(&batch_flow, registry).await
265        }
266        CreditRegistrationPhase::Import => import::run(&batch_flow, registry).await,
267        CreditRegistrationPhase::Verify => verify::run(&batch_flow, registry).await,
268        CreditRegistrationPhase::LegacyMirror => legacy_mirror::run(pool, scope).await,
269        CreditRegistrationPhase::StudentNotifications => {
270            student_notifications::run(pool, scope, ctx.base_url).await
271        }
272        CreditRegistrationPhase::EnrolmentDiscovery => {
273            enrolment_discovery::run(pool, scope, ctx.is_account_linking_enabled, registry).await
274        }
275        CreditRegistrationPhase::LinkEmails => link_emails::run(pool, scope, ctx.base_url).await,
276        CreditRegistrationPhase::ConfigValidation => {
277            config_validation::run(pool, scope, registry).await
278        }
279        CreditRegistrationPhase::RetentionSweep => retention_sweep::run(pool).await,
280        CreditRegistrationPhase::LedgerSnapshot => ledger_snapshot::run(pool).await,
281    }
282}
283
284#[cfg(test)]
285mod tests {
286    use super::*;
287
288    #[test]
289    fn the_audit_name_says_who_ran_the_phase() {
290        for caller in [WorkerProcess::CreditRegistrar.as_str(), "run-tick"] {
291            for phase in CreditRegistrationPhase::ALL {
292                let name = worker_name(caller, phase);
293                assert!(name.starts_with(caller));
294                assert!(name.ends_with(phase.as_str()));
295                assert!(name.len() <= 64, "{name}");
296            }
297        }
298    }
299}