headless_lms_utils/
periodic_worker.rs1use std::time::Duration;
5
6use tokio_util::sync::CancellationToken;
7
8pub struct PeriodicWorkerConfig<'a> {
11 pub tick_interval: Duration,
12 pub still_running: Option<StillRunningLog<'a>>,
13 pub delay_missed_ticks: bool,
16}
17
18pub struct StillRunningLog<'a> {
20 pub every: u32,
21 pub message: &'a str,
22 pub initial_ticks: u32,
25}
26
27pub 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
36pub 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}