headless_lms_credit_registration/runtime/
worker_loop.rs1use 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
35const TICK_INTERVAL: Duration = Duration::from_secs(10);
38
39const STILL_RUNNING_MESSAGE_TICKS: u32 = 60;
41
42pub 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 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 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 still_running: None,
107 delay_missed_ticks: true,
110 },
111 shutdown,
112 async || {
113 trace!(phase = phase.as_str(), "Checking whether phase is due");
114 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
132async 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
173async fn run_due_phase(
175 ctx: &PhaseContext<'_>,
176 phase: CreditRegistrationPhase,
177 state: &CreditRegistrationPhaseState,
178) -> CreditRegistrationResult<()> {
179 let mut conn = ctx.pool.acquire().await?;
180 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 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 trace!(
218 phase = phase.as_str(),
219 duration_ms, "Credit registration phase run found nothing to do"
220 );
221 }
222 PhaseTick::Skipped(reason) => log_skip_if_changed(phase, reason),
225 PhaseTick::ScopeNotSupported => {}
226 }
227 Ok(())
228}
229
230static 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
272fn 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 #[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 #[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}