Skip to main content

headless_lms_utils/
periodic_worker.rs

1//! The tick-interval scaffold shared by every background worker that polls the database on a fixed
2//! period: `regrader`, `chatbot_syncer` and the credit registration phase runners.
3
4use std::time::Duration;
5
6use tokio_util::sync::CancellationToken;
7
8/// How a worker's ticking loop should behave, so the loop itself carries none of that per-worker
9/// detail.
10pub struct PeriodicWorkerConfig<'a> {
11    pub tick_interval: Duration,
12    pub still_running: Option<StillRunningLog<'a>>,
13    /// `true` pushes a slow iteration's next tick out instead of firing it immediately
14    /// (`tokio::time::MissedTickBehavior::Delay`); `false` keeps tokio's default (`Burst`).
15    pub delay_missed_ticks: bool,
16}
17
18/// The "still running" heartbeat line a worker logs every `every` ticks.
19pub struct StillRunningLog<'a> {
20    pub every: u32,
21    pub message: &'a str,
22    /// Starting value of the tick counter, so a worker that wants its first line sooner than `every`
23    /// ticks can seed it.
24    pub initial_ticks: u32,
25}
26
27/// Runs `body` on `config.tick_interval` forever, logging the still-running line if there is one. A
28/// `body` that returns `Err` stops the loop and becomes this function's return value.
29pub async fn run_periodic_worker(
30    config: PeriodicWorkerConfig<'_>,
31    body: impl AsyncFnMut() -> anyhow::Result<()>,
32) -> anyhow::Result<()> {
33    run_periodic_worker_until(config, &CancellationToken::new(), body).await
34}
35
36/// [`run_periodic_worker`] that returns `Ok` once `shutdown` is cancelled. An iteration already
37/// running is left to finish; only the wait for the next tick is cut short.
38pub async fn run_periodic_worker_until<E>(
39    config: PeriodicWorkerConfig<'_>,
40    shutdown: &CancellationToken,
41    mut body: impl AsyncFnMut() -> Result<(), E>,
42) -> Result<(), E> {
43    let mut interval = tokio::time::interval(config.tick_interval);
44    if config.delay_missed_ticks {
45        interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
46    }
47    let mut ticks = config
48        .still_running
49        .as_ref()
50        .map_or(0, |still_running| still_running.initial_ticks);
51    loop {
52        tokio::select! {
53            biased;
54            _ = shutdown.cancelled() => return Ok(()),
55            _ = interval.tick() => {}
56        }
57        if let Some(still_running) = &config.still_running {
58            ticks += 1;
59            if ticks >= still_running.every {
60                ticks = 0;
61                info!("{}", still_running.message);
62            }
63        }
64        body().await?;
65    }
66}