headless_lms_credit_registration/runtime/
dispatch.rs1use 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#[derive(Debug, Clone, PartialEq)]
29pub enum PhaseTick {
30 Ran(PhaseRunOutcome),
31 Skipped(PhaseSkipReason),
33 ScopeNotSupported,
35}
36
37#[derive(Debug, Clone, Copy, PartialEq, Eq)]
38pub enum PhaseSkipReason {
39 Paused,
40 CircuitBreakerOpen,
41 AccountLinkingDisabled,
42 SisuDayGap,
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub enum Runner<'a> {
49 Worker(WorkerProcess),
52 Other(&'a str),
55}
56
57impl Runner<'_> {
58 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
74pub struct PhaseContext<'a> {
76 pub pool: &'a PgPool,
77 pub suotar_client: &'a SuotarClient,
78 pub test_mode: bool,
80 pub runner: Runner<'a>,
81 pub base_url: &'a str,
83 pub is_account_linking_enabled: bool,
85 pub shutdown: Option<&'a CancellationToken>,
87}
88
89impl<'a> PhaseContext<'a> {
90 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
110pub(super) fn worker_name(caller: &str, phase: CreditRegistrationPhase) -> String {
112 format!("{caller}/{}", phase.as_str())
113}
114
115#[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 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 let bookkeeping = scope.is_unscoped();
136 let is_own_worker = ctx.runner.owning_process() == Some(phase.spec().process);
137 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 if phase.spec().is_account_linking_only && !ctx.is_account_linking_enabled {
146 return Ok(PhaseTick::Skipped(PhaseSkipReason::AccountLinkingDisabled));
147 }
148 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 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
240async 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 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}