85 lines
2.5 KiB
Rust
85 lines
2.5 KiB
Rust
//! Cron scheduler: decides which scheduled job kinds are due and runs them.
|
|
use std::str::FromStr;
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use chrono::{DateTime, Utc};
|
|
use cron::Schedule;
|
|
use domain::jobs::JobKind;
|
|
use domain::DomainError;
|
|
|
|
use crate::{JobRunner, SettingsService};
|
|
|
|
/// Validate a 6-field cron expression (seconds first).
|
|
pub fn validate_cron(expr: &str) -> Result<(), DomainError> {
|
|
Schedule::from_str(expr)
|
|
.map(|_| ())
|
|
.map_err(|e| DomainError::Validation(format!("invalid cron expression: {e}")))
|
|
}
|
|
|
|
/// Next fire time strictly after `after`.
|
|
pub fn next_fire(expr: &str, after: DateTime<Utc>) -> Option<DateTime<Utc>> {
|
|
Schedule::from_str(expr).ok()?.after(&after).next()
|
|
}
|
|
|
|
/// A job is due if a fire time exists in `(last_run, now]`. With no last run, it is due
|
|
/// if a fire time falls within the last `grace` seconds, so a fresh start does not
|
|
/// immediately run every job.
|
|
pub fn is_due(
|
|
expr: &str,
|
|
last_run: Option<DateTime<Utc>>,
|
|
now: DateTime<Utc>,
|
|
grace_secs: i64,
|
|
) -> bool {
|
|
let since = last_run.unwrap_or(now - chrono::Duration::seconds(grace_secs));
|
|
next_fire(expr, since).is_some_and(|t| t <= now)
|
|
}
|
|
|
|
pub struct Scheduler {
|
|
runner: Arc<JobRunner>,
|
|
settings: Arc<SettingsService>,
|
|
tick: Duration,
|
|
}
|
|
|
|
impl Scheduler {
|
|
pub fn new(runner: Arc<JobRunner>, settings: Arc<SettingsService>) -> Self {
|
|
Self {
|
|
runner,
|
|
settings,
|
|
tick: Duration::from_secs(30),
|
|
}
|
|
}
|
|
|
|
/// One pass: start every registered kind that is due. Returns the kinds started.
|
|
pub async fn tick_once(&self, now: DateTime<Utc>) -> Vec<JobKind> {
|
|
let mut started = Vec::new();
|
|
for kind in self.runner.kinds() {
|
|
let Ok(Some(cron)) = self.settings.schedule(kind).await else {
|
|
continue;
|
|
};
|
|
let last = self
|
|
.runner
|
|
.last_finished(kind)
|
|
.await
|
|
.ok()
|
|
.flatten()
|
|
.map(|r| r.started_at);
|
|
if is_due(&cron, last, now, self.tick.as_secs() as i64 * 2)
|
|
&& self.runner.start(kind, None, "scheduler").await.is_ok()
|
|
{
|
|
started.push(kind);
|
|
}
|
|
}
|
|
started
|
|
}
|
|
|
|
/// Run forever; meant to be spawned on the runtime.
|
|
pub async fn run(self) {
|
|
let mut interval = tokio::time::interval(self.tick);
|
|
loop {
|
|
interval.tick().await;
|
|
self.tick_once(Utc::now()).await;
|
|
}
|
|
}
|
|
}
|