use std::collections::BTreeMap;
use std::sync::Mutex;
use std::time::Duration;
use jiff::Timestamp;
use layover_core::config::Config;
use layover_core::pipeline::{PipelineName, Schedule};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Skipped {
StillWorking {
pipeline: PipelineName,
},
}
impl std::fmt::Display for Skipped {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::StillWorking { pipeline } => write!(
f,
"`{pipeline}` was due but its previous run has not finished; \
set `overlap = \"allow\"` if two at once is safe"
),
}
}
}
#[derive(Debug, Default, PartialEq, Eq)]
pub struct Due {
pub fire: Vec<PipelineName>,
pub skipped: Vec<Skipped>,
}
#[derive(Debug)]
pub struct Clock {
next: Mutex<BTreeMap<PipelineName, Timestamp>>,
}
impl Clock {
#[must_use]
pub fn new(config: &Config, now: Timestamp) -> Self {
let mut next = BTreeMap::new();
for (name, pipeline) in &config.pipelines {
if let Some(schedule) = pipeline.trigger.schedule()
&& let Some(at) = next_after(schedule, now)
{
next.insert(name.clone(), at);
}
}
Self {
next: Mutex::new(next),
}
}
pub fn tick(
&self,
config: &Config,
now: Timestamp,
working: impl Fn(&PipelineName) -> bool,
) -> Due {
let Ok(mut next) = self.next.lock() else {
return Due::default();
};
let mut due = Due::default();
for (name, at) in next.iter_mut() {
if *at > now {
continue;
}
if let Some(schedule) = config
.pipelines
.get(name)
.and_then(|p| p.trigger.schedule())
&& let Some(following) = next_after(schedule, now)
{
*at = following;
}
let overlaps = config
.pipelines
.get(name)
.is_some_and(layover_core::pipeline::Pipeline::allows_overlap);
if overlaps || !working(name) {
due.fire.push(name.clone());
} else {
due.skipped.push(Skipped::StillWorking {
pipeline: name.clone(),
});
}
}
due
}
#[must_use]
pub fn next_due(&self, pipeline: &PipelineName) -> Option<Timestamp> {
self.next.lock().ok()?.get(pipeline).copied()
}
#[must_use]
pub fn until_next(&self, now: Timestamp) -> Option<Duration> {
let next = self.next.lock().ok()?;
let soonest = next.values().min()?;
Some(
soonest
.duration_since(now)
.try_into()
.unwrap_or(Duration::ZERO),
)
}
}
fn next_after(schedule: &Schedule, now: Timestamp) -> Option<Timestamp> {
schedule.next_after(now)
}
#[cfg(test)]
mod tests {
use super::*;
use jiff::ToSpan as _;
fn factory(pipelines: &str) -> Config {
let text = format!(
r#"
[layover]
work_dir = "work"
[defaults]
runner = "shell"
[runners.shell]
command = ["echo"]
[agents.worker]
prompt = "work"
entry = true
{pipelines}
"#
);
toml::from_str(&text).expect("the fixture factory parses")
}
fn name(text: &str) -> PipelineName {
PipelineName::new(text)
}
fn never_working(_: &PipelineName) -> bool {
false
}
#[test]
fn nothing_fires_at_startup() {
let config = factory(
r#"
[pipelines.sweep]
entry = "worker"
trigger = { every = "1h" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let due = clock.tick(&config, now, never_working);
assert!(due.fire.is_empty(), "{:?}", due.fire);
assert!(due.skipped.is_empty());
}
#[test]
fn an_interval_pipeline_fires_once_its_interval_has_passed() {
let config = factory(
r#"
[pipelines.sweep]
entry = "worker"
trigger = { every = "1h" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let later = now.checked_add(61_i32.minutes()).expect("in range");
let due = clock.tick(&config, later, never_working);
assert_eq!(due.fire, [name("sweep")]);
}
#[test]
fn a_tower_asleep_for_hours_fires_once_on_waking_not_once_per_missed_tick() {
let config = factory(
r#"
[pipelines.sweep]
entry = "worker"
trigger = { every = "1h" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let much_later = now.checked_add(6_i32.hours()).expect("in range");
let due = clock.tick(&config, much_later, never_working);
assert_eq!(due.fire, [name("sweep")], "one firing, not six");
}
#[test]
fn a_pipeline_that_is_still_working_is_skipped_and_says_why() {
let config = factory(
r#"
[pipelines.sweep]
entry = "worker"
trigger = { every = "1h" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let later = now.checked_add(61_i32.minutes()).expect("in range");
let due = clock.tick(&config, later, |_| true);
assert!(due.fire.is_empty());
assert_eq!(
due.skipped,
[Skipped::StillWorking {
pipeline: name("sweep")
}]
);
assert!(
due.skipped[0].to_string().contains("overlap"),
"a skip should say how to opt out of it: {}",
due.skipped[0]
);
}
#[test]
fn overlap_allow_starts_a_second_copy_deliberately() {
let config = factory(
r#"
[pipelines.sweep]
entry = "worker"
trigger = { every = "1h" }
overlap = "allow"
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let later = now.checked_add(61_i32.minutes()).expect("in range");
let due = clock.tick(&config, later, |_| true);
assert_eq!(due.fire, [name("sweep")]);
assert!(due.skipped.is_empty());
}
#[test]
fn a_skipped_tick_still_advances_the_clock() {
let config = factory(
r#"
[pipelines.sweep]
entry = "worker"
trigger = { every = "1h" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let later = now.checked_add(61_i32.minutes()).expect("in range");
clock.tick(&config, later, |_| true);
let again = clock.tick(&config, later, never_working);
assert!(
again.fire.is_empty(),
"the same tick must not fire twice: {:?}",
again.fire
);
}
#[test]
fn a_manual_pipeline_is_never_due() {
let config = factory(
r#"
[pipelines.onbehalf]
entry = "worker"
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let far = now.checked_add(30_i32.hours()).expect("in range");
assert!(clock.tick(&config, far, never_working).fire.is_empty());
assert!(clock.until_next(now).is_none(), "nothing to wait for");
}
#[test]
fn a_cron_pipeline_is_due_at_its_next_matching_minute() {
let config = factory(
r#"
[pipelines.digest]
entry = "worker"
trigger = { cron = "* * * * *" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let next = clock.next_due(&name("digest")).expect("scheduled");
assert!(next > now, "the next firing is in the future");
assert!(
next.duration_since(now).as_secs() <= 60,
"an every-minute cron should be due within the minute"
);
}
#[test]
fn the_wait_is_until_the_soonest_pipeline_not_the_first_one_declared() {
let config = factory(
r#"
[pipelines.slow]
entry = "worker"
trigger = { every = "6h" }
[pipelines.quick]
entry = "worker"
trigger = { every = "5m" }
"#,
);
let now = Timestamp::now();
let clock = Clock::new(&config, now);
let wait = clock.until_next(now).expect("something is scheduled");
assert!(
wait <= Duration::from_mins(5),
"should wait for `quick`, not `slow`: {wait:?}"
);
}
}