(job: ScheduledJob, sink: Arc<dyn ScheduleSink>, cancel: CancellationToken)
| 95 | } |
| 96 | |
| 97 | async fn run_job(job: ScheduledJob, sink: Arc<dyn ScheduleSink>, cancel: CancellationToken) { |
| 98 | loop { |
| 99 | let now = Utc::now(); |
| 100 | let Some(next) = job.next_fire_after(now) else { |
| 101 | return; // never fires again |
| 102 | }; |
| 103 | let wait = (next - now) |
| 104 | .to_std() |
| 105 | .unwrap_or(std::time::Duration::from_secs(0)); |
| 106 | tokio::select! { |
| 107 | _ = tokio::time::sleep(wait) => sink.fire(&job.spec).await, |
| 108 | _ = cancel.cancelled() => return, |
| 109 | } |
| 110 | } |
| 111 | } |
| 112 | |
| 113 | #[cfg(test)] |
| 114 | mod tests { |
no test coverage detected