diff --git a/CHANGELOG.md b/CHANGELOG.md index 3a2e46d..6e1271d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- **`crates/tinyflows-schedule`** — the schedule model (`Schedule`, + `ActiveHours`, the job/run records) and its pure logic: cron-expression + normalisation, time-zone and active-window aware next-run computation, + validation and the minimum-cadence check. Extracted from OpenHuman's `cron` + domain; the serde shapes are byte-identical and pinned by literal-JSON + fixtures. No runtime, store or config. - **`tinyflows_catalog::graph_hash::compute_graph_hash`** — the content pin over a graph plus its `require_approval` flag that a parked run records and a resume re-checks. Extracted from OpenHuman; the digest is byte-identical to diff --git a/CLAUDE.md b/CLAUDE.md index 9e5fda7..863b93e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -26,6 +26,7 @@ uses, so a reader who knows one knows them all. | `crates/tinyflows-sqlite` | SQLite implementations of the catalog and of the engine's checkpoint store, plus the JSON draft store. Every entry point takes a directory, never a host config type. | | `crates/tinyflows-copilot` | The authoring copilot's *words*: the `workflow_builder` / `flow_discovery` standing archetypes, and the turn brief that opens a builder turn. Names no tool trait, no agent registry, no model client. | | `crates/tinyflows-adaptive` | The adaptive loop over the engine: select or author a workflow, run it, judge it, learn. | +| `crates/tinyflows-schedule` | Schedule model and pure next-run logic (cron / interval / one-shot, time zones, active hours, cadence floor). No runtime, store or config. | **Where a new thing goes.** Ask what it depends on, not what it is about. If it needs storage, it is not `tinyflows-catalog`. If it needs a tool trait or a model diff --git a/Cargo.lock b/Cargo.lock index d59dcd8..1eaf8f7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -235,6 +235,16 @@ dependencies = [ "windows-link", ] +[[package]] +name = "chrono-tz" +version = "0.10.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6139a8597ed92cf816dfb33f5dd6cf0bb93a6adc938f11039f371bc5bcd26c3" +dependencies = [ + "chrono", + "phf", +] + [[package]] name = "cmake" version = "0.1.58" @@ -345,6 +355,17 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" +[[package]] +name = "cron" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6f8c3e73077b4b4a6ab1ea5047c37c57aee77657bc8ecd6f29b0af082d0b0c07" +dependencies = [ + "chrono", + "nom", + "once_cell", +] + [[package]] name = "crossbeam-channel" version = "0.5.16" @@ -1384,6 +1405,12 @@ version = "2.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "88904434abc2901f197fe8cc55f0445e7ded921dba5911dad2e2b39b48e663c4" +[[package]] +name = "minimal-lexical" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68354c5c6bd36d73ff3feceb05efa59b6acb7626617f4962be322a825e61f79a" + [[package]] name = "miniz_oxide" version = "0.8.9" @@ -1514,6 +1541,16 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "27b02d87554356db9e9a873add8782d4ea6e3e58ea071a9adb9a2e8ddb884a8b" +[[package]] +name = "nom" +version = "7.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d273983c5a657a70a3e8f2a01329822f3b8c8172b73826411a55751e404a0a4a" +dependencies = [ + "memchr", + "minimal-lexical", +] + [[package]] name = "num-bigint" version = "0.4.6" @@ -1602,6 +1639,24 @@ version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" +[[package]] +name = "phf" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "913273894cec178f401a31ec4b656318d95473527be05c0752cc41cdc32be8b7" +dependencies = [ + "phf_shared", +] + +[[package]] +name = "phf_shared" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06005508882fb681fd97892ecff4b7fd0fee13ef1aa569f8695dae7ab9099981" +dependencies = [ + "siphasher", +] + [[package]] name = "pin-project-lite" version = "0.2.17" @@ -2293,6 +2348,12 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e3a9fe34e3e7a50316060351f37187a3f546bce95496156754b601a5fa71b76e" +[[package]] +name = "siphasher" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "33f4fe9184a62d842c9ef383018f3306d8ba224fd9d836f56d7288308847c256" + [[package]] name = "slab" version = "0.4.12" @@ -2571,6 +2632,19 @@ dependencies = [ "serde_json", ] +[[package]] +name = "tinyflows-schedule" +version = "0.1.0" +dependencies = [ + "anyhow", + "chrono", + "chrono-tz", + "cron", + "serde", + "serde_json", + "tracing", +] + [[package]] name = "tinyflows-sqlite" version = "0.1.0" diff --git a/crates/tinyflows-schedule/Cargo.toml b/crates/tinyflows-schedule/Cargo.toml new file mode 100644 index 0000000..54c6251 --- /dev/null +++ b/crates/tinyflows-schedule/Cargo.toml @@ -0,0 +1,28 @@ +# Pure scheduling logic for timed jobs: the `Schedule` / `ActiveHours` model +# (with its persisted wire shape), cron-expression normalisation, time-zone and +# active-window aware next-run computation, and the minimum-cadence check. +# +# Deliberately no scheduler runtime, store, clock or config here: a host owns +# when to poll, where jobs are kept and what a run does. This crate only +# answers "given this schedule and this instant, when next?". +[package] +name = "tinyflows-schedule" +version = "0.1.0" +description = "Cron/interval/one-shot schedule model and next-run computation for tinyflows hosts" +edition.workspace = true +rust-version.workspace = true +license.workspace = true +repository.workspace = true +homepage.workspace = true +readme = "../../README.md" + +[dependencies] +serde = { workspace = true } +chrono = { workspace = true } +anyhow = "1.0" +cron = "0.12" +chrono-tz = "0.10" +tracing = { workspace = true } + +[dev-dependencies] +serde_json = { workspace = true } diff --git a/crates/tinyflows-schedule/src/lib.rs b/crates/tinyflows-schedule/src/lib.rs new file mode 100644 index 0000000..eee04f6 --- /dev/null +++ b/crates/tinyflows-schedule/src/lib.rs @@ -0,0 +1,24 @@ +//! Scheduling model and next-run computation for timed jobs. +//! +//! * [`types`] — [`Schedule`], [`ActiveHours`] and the job/run records a host +//! persists. Their serde shapes are a wire contract (a host's job store and +//! RPC surface carry them), pinned by `wire_tests`. +//! * [`schedule`] — next-run computation, validation, cron-expression +//! normalisation and the minimum-cadence check. +//! +//! No runtime, storage or configuration lives here; a host supplies those. + +pub mod schedule; +pub mod types; + +pub use schedule::{ + MIN_AGENT_JOB_INTERVAL, TooFrequent, next_run_for_schedule, normalize_expression, + runs_closer_than, schedule_cron_expression, validate_agent_schedule, validate_schedule, +}; +pub use types::{ + ActiveHours, CronJob, CronJobPatch, CronRun, DeliveryConfig, JobType, Schedule, SessionTarget, +}; + +#[cfg(test)] +#[path = "wire_tests.rs"] +mod wire_tests; diff --git a/crates/tinyflows-schedule/src/schedule.rs b/crates/tinyflows-schedule/src/schedule.rs new file mode 100644 index 0000000..e8033f8 --- /dev/null +++ b/crates/tinyflows-schedule/src/schedule.rs @@ -0,0 +1,374 @@ +use crate::types::{ActiveHours, Schedule}; +use anyhow::{Context, Result}; +use chrono::{DateTime, Duration as ChronoDuration, NaiveTime, Timelike, Utc}; +use cron::Schedule as CronExprSchedule; +use std::fmt; +use std::str::FromStr; + +/// The closest together two runs of an *agent* job may be scheduled. Every +/// run is a full inference turn, so a tighter cadence is almost always a +/// misconfiguration that bills accordingly. Shell and flow jobs are not +/// subject to it. Enforced by [`validate_agent_schedule`] when an agent job is +/// created or its schedule changed; the scheduler warns about rows that +/// predate the rule. +pub const MIN_AGENT_JOB_INTERVAL: ChronoDuration = ChronoDuration::minutes(5); + +/// Upper bound on cron candidates walked while looking for an occurrence that +/// falls inside `active_hours`. `next_run_for_schedule` gets a fresh budget +/// per call; a `runs_closer_than` scan shares one across all of its steps. +const ACTIVE_WINDOW_CANDIDATE_LIMIT: usize = 600_000; + +/// How many consecutive occurrences [`runs_closer_than`] walks before it +/// concludes a cron schedule keeps its distance. A full leap year at the +/// minimum allowed cadence has at most 105,408 runs (288 per day), so this +/// covers annual timezone transitions regardless of where the scan starts. +const RUN_GAP_SCAN_OCCURRENCES: usize = 105_500; + +pub fn next_run_for_schedule(schedule: &Schedule, from: DateTime) -> Result> { + match schedule { + Schedule::Cron { + expr, + tz, + active_hours, + } => { + let plan = CronPlan::parse(expr, tz.as_deref(), active_hours.as_ref())?; + let mut budget = ACTIVE_WINDOW_CANDIDATE_LIMIT; + plan.next_after(from, &mut budget) + } + Schedule::At { at } => Ok(*at), + Schedule::Every { every_ms } => { + if *every_ms == 0 { + anyhow::bail!("Invalid schedule: every_ms must be > 0"); + } + let ms = i64::try_from(*every_ms).context("every_ms is too large")?; + let delta = ChronoDuration::milliseconds(ms); + from.checked_add_signed(delta) + .ok_or_else(|| anyhow::anyhow!("every_ms overflowed DateTime")) + } + } +} + +pub fn validate_schedule(schedule: &Schedule, now: DateTime) -> Result<()> { + match schedule { + Schedule::Cron { + expr, + tz, + active_hours, + } => { + let _ = normalize_expression(expr)?; + if let Some(active) = active_hours { + let _ = ActiveWindow::parse(active)?; + } + let _ = ScheduleTimeZone::parse(tz.as_deref())?; + let _ = next_run_for_schedule(schedule, now)?; + Ok(()) + } + Schedule::At { at } => { + if *at <= now { + anyhow::bail!("Invalid schedule: 'at' must be in the future"); + } + Ok(()) + } + Schedule::Every { every_ms } => { + if *every_ms == 0 { + anyhow::bail!("Invalid schedule: every_ms must be > 0"); + } + // Validate against the same conversion and checked arithmetic used + // when computing the next run, so persisted schedules are usable. + let _ = next_run_for_schedule(schedule, now)?; + Ok(()) + } + } +} + +/// [`validate_schedule`] plus the agent-only floor: an agent job may not run +/// closer together than [`MIN_AGENT_JOB_INTERVAL`]. The error names the two +/// runs (or the fixed interval) that break the rule, so the caller — an agent +/// using `cron_add`, or the settings form — can say exactly what to change. +pub fn validate_agent_schedule(schedule: &Schedule, now: DateTime) -> Result<()> { + validate_schedule(schedule, now)?; + if let Some(too_frequent) = runs_closer_than(schedule, now, MIN_AGENT_JOB_INTERVAL) { + anyhow::bail!( + "Invalid schedule: agent jobs must run at least {} apart, but this schedule {too_frequent}", + describe_gap(MIN_AGENT_JOB_INTERVAL) + ); + } + Ok(()) +} + +/// Why a schedule runs more often than a threshold allows. Carries the +/// evidence, not only the verdict, so a log line or an error can quote it. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TooFrequent { + /// A fixed `every_ms` interval shorter than the threshold. + FixedInterval(ChronoDuration), + /// Two consecutive cron occurrences closer together than the threshold. + ConsecutiveRuns { + first: DateTime, + second: DateTime, + }, +} + +impl TooFrequent { + /// The offending gap. + pub fn gap(&self) -> ChronoDuration { + match self { + Self::FixedInterval(gap) => *gap, + Self::ConsecutiveRuns { first, second } => *second - *first, + } + } +} + +impl fmt::Display for TooFrequent { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::FixedInterval(gap) => write!(f, "fires every {}", describe_gap(*gap)), + Self::ConsecutiveRuns { first, second } => write!( + f, + "fires at {} and again at {}, {} apart", + first.format("%Y-%m-%d %H:%M:%S UTC"), + second.format("%Y-%m-%d %H:%M:%S UTC"), + describe_gap(*second - *first) + ), + } + } +} + +/// The two consecutive runs of `schedule` after `from` that are closest +/// together, if that gap is under `min_gap`. +/// +/// Consecutive occurrences are walked in order and the smallest gap is kept, +/// so an irregular expression such as `1,3,4,30 * * * *` is judged by its +/// :03 → :04 pair — not by the first pair under the floor (:01 → :03) and not +/// by whichever pair happens to follow `from`. The wrap-around gap counts too: +/// `*/7 * * * *` fires at :56 and then at :00, four minutes apart, and is +/// reported as such. Ties keep the earliest pair. +/// +/// The walk is bounded ([`RUN_GAP_SCAN_OCCURRENCES`] runs, one shared +/// [`ACTIVE_WINDOW_CANDIDATE_LIMIT`] budget), so a sparse or window-restricted +/// expression stays cheap; if the budget runs out mid-walk, the closest pair +/// seen so far is still reported. An expression that cannot be parsed, or that +/// has no second occurrence, is not evidence of anything and yields `None`; +/// [`validate_schedule`] is where a bad expression gets rejected. +pub fn runs_closer_than( + schedule: &Schedule, + from: DateTime, + min_gap: ChronoDuration, +) -> Option { + match schedule { + Schedule::At { .. } => None, + Schedule::Every { every_ms } => { + let gap = ChronoDuration::try_milliseconds(i64::try_from(*every_ms).ok()?)?; + (gap < min_gap).then_some(TooFrequent::FixedInterval(gap)) + } + Schedule::Cron { + expr, + tz, + active_hours, + } => { + let plan = CronPlan::parse(expr, tz.as_deref(), active_hours.as_ref()).ok()?; + let mut budget = ACTIVE_WINDOW_CANDIDATE_LIMIT; + let mut previous = plan.next_after(from, &mut budget).ok()?; + let mut closest: Option = None; + for _ in 1..RUN_GAP_SCAN_OCCURRENCES { + // Running out of budget (or of occurrences) ends the walk but + // does not discard a pair already found. + let Ok(next) = plan.next_after(previous, &mut budget) else { + break; + }; + let gap = next - previous; + if gap < min_gap && closest.is_none_or(|seen| gap < seen.gap()) { + closest = Some(TooFrequent::ConsecutiveRuns { + first: previous, + second: next, + }); + } + previous = next; + } + closest + } + } +} + +/// `4 minutes`, `30 seconds`, `1 minute 30 seconds` — whole seconds only. +fn describe_gap(gap: ChronoDuration) -> String { + fn count(n: i64, unit: &str) -> String { + if n == 1 { + format!("{n} {unit}") + } else { + format!("{n} {unit}s") + } + } + let seconds = gap.num_seconds(); + match (seconds / 60, seconds % 60) { + (0, s) => count(s, "second"), + (m, 0) => count(m, "minute"), + (m, s) => format!("{} {}", count(m, "minute"), count(s, "second")), + } +} + +pub fn schedule_cron_expression(schedule: &Schedule) -> Option { + match schedule { + Schedule::Cron { expr, .. } => Some(expr.clone()), + _ => None, + } +} + +/// A [`Schedule::Cron`] parsed once, so walking many occurrences does not pay +/// for the expression, timezone and active-window parsing on every step. +struct CronPlan<'a> { + expr: &'a str, + cron: CronExprSchedule, + timezone: ScheduleTimeZone, + active_window: Option, +} + +impl<'a> CronPlan<'a> { + fn parse(expr: &'a str, tz: Option<&str>, active_hours: Option<&ActiveHours>) -> Result { + let normalized = normalize_expression(expr)?; + let cron = CronExprSchedule::from_str(&normalized) + .with_context(|| format!("Invalid cron expression: {expr}"))?; + let timezone = ScheduleTimeZone::parse(tz)?; + let active_window = active_hours.map(ActiveWindow::parse).transpose()?; + Ok(Self { + expr, + cron, + timezone, + active_window, + }) + } + + /// The first occurrence strictly after `from` that falls inside the active + /// window, spending at most `budget` cron candidates to find it. + fn next_after(&self, from: DateTime, budget: &mut usize) -> Result> { + let mut current_from = from; + while *budget > 0 { + *budget -= 1; + let next_utc = self + .timezone + .next_after(&self.cron, current_from, self.expr)?; + let Some(active) = &self.active_window else { + return Ok(next_utc); + }; + if active.contains(self.timezone.local_time_of_day(next_utc)) { + return Ok(next_utc); + } + tracing::debug!( + "[cron] next_run candidate {} outside active window {}–{}, advancing", + next_utc, + active.start, + active.end + ); + current_from = next_utc; + } + tracing::warn!( + "[cron] no occurrence found within active_hours for expr={} after 100,000 candidates", + self.expr + ); + anyhow::bail!("No future occurrence found within active hours after 100,000 attempts") + } +} + +#[derive(Debug, Clone, Copy)] +enum ScheduleTimeZone { + Local, + Named(chrono_tz::Tz), +} + +impl ScheduleTimeZone { + fn parse(tz: Option<&str>) -> Result { + match tz { + Some(tz_name) => chrono_tz::Tz::from_str(tz_name) + .map(Self::Named) + .with_context(|| format!("Invalid IANA timezone: {tz_name}")), + None => Ok(Self::Local), + } + } + + fn next_after( + self, + cron: &CronExprSchedule, + from: DateTime, + expr: &str, + ) -> Result> { + match self { + Self::Named(timezone) => { + let localized_from = from.with_timezone(&timezone); + let next_local = cron.after(&localized_from).next().ok_or_else(|| { + anyhow::anyhow!("No future occurrence for expression: {expr}") + })?; + Ok(next_local.with_timezone(&Utc)) + } + Self::Local => { + let localized_from = from.with_timezone(&chrono::Local); + let next_local = cron.after(&localized_from).next().ok_or_else(|| { + anyhow::anyhow!("No future occurrence for expression: {expr}") + })?; + Ok(next_local.with_timezone(&Utc)) + } + } + } + + fn local_time_of_day(self, time: DateTime) -> NaiveTime { + match self { + Self::Named(timezone) => { + let localized = time.with_timezone(&timezone); + NaiveTime::from_hms_opt(localized.hour(), localized.minute(), 0) + .expect("hour() and minute() from a valid DateTime are always in-range") + } + Self::Local => { + let localized = time.with_timezone(&chrono::Local); + NaiveTime::from_hms_opt(localized.hour(), localized.minute(), 0) + .expect("hour() and minute() from a valid DateTime are always in-range") + } + } + } +} + +#[derive(Debug, Clone, Copy)] +struct ActiveWindow { + start: NaiveTime, + end: NaiveTime, +} + +impl ActiveWindow { + fn parse(active: &ActiveHours) -> Result { + let start = NaiveTime::parse_from_str(&active.start, "%H:%M") + .with_context(|| format!("Invalid active_hours.start: {}", active.start))?; + let end = NaiveTime::parse_from_str(&active.end, "%H:%M") + .with_context(|| format!("Invalid active_hours.end: {}", active.end))?; + Ok(Self { start, end }) + } + + fn contains(self, time: NaiveTime) -> bool { + if self.start <= self.end { + time >= self.start && time <= self.end + } else { + // Window spans midnight (e.g. 22:00 to 06:00). + time >= self.start || time <= self.end + } + } +} + +pub fn normalize_expression(expression: &str) -> Result { + let expression = expression.trim(); + let field_count = expression.split_whitespace().count(); + + match field_count { + // standard crontab syntax: minute hour day month weekday + 5 => Ok(format!("0 {expression}")), + // crate-native syntax includes seconds (+ optional year) + 6 | 7 => Ok(expression.to_string()), + _ => anyhow::bail!( + "Invalid cron expression: {expression} (expected 5, 6, or 7 fields, got {field_count})" + ), + } +} + +#[cfg(test)] +#[path = "schedule_tests.rs"] +mod tests; + +#[cfg(test)] +#[path = "schedule_gap_tests.rs"] +mod gap_tests; diff --git a/crates/tinyflows-schedule/src/schedule_gap_tests.rs b/crates/tinyflows-schedule/src/schedule_gap_tests.rs new file mode 100644 index 0000000..894d729 --- /dev/null +++ b/crates/tinyflows-schedule/src/schedule_gap_tests.rs @@ -0,0 +1,249 @@ +use super::*; +use chrono::TimeZone; + +// ── runs_closer_than / validate_agent_schedule (#6158) ────────── + +fn utc_cron(expr: &str) -> Schedule { + Schedule::Cron { + expr: expr.into(), + tz: Some("UTC".into()), + active_hours: None, + } +} + +#[test] +fn runs_closer_than_reports_the_shortest_gap_of_an_irregular_expression() { + // From :03 the pair that follows is :30 → 1:01 — 31 minutes apart, which + // is what a "the pair after now" check would have judged the schedule by. + // The :01 → :02 pair one minute apart is what counts. + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, 3, 0).unwrap(); + let hit = runs_closer_than(&utc_cron("1,2,30 * * * *"), from, MIN_AGENT_JOB_INTERVAL) + .expect("the one-minute gap must be found"); + assert_eq!( + hit, + TooFrequent::ConsecutiveRuns { + first: Utc.with_ymd_and_hms(2026, 3, 2, 11, 1, 0).unwrap(), + second: Utc.with_ymd_and_hms(2026, 3, 2, 11, 2, 0).unwrap(), + } + ); + assert_eq!(hit.gap(), ChronoDuration::minutes(1)); +} + +#[test] +fn runs_closer_than_does_not_depend_on_the_instant_it_starts_from() { + for minute in [0, 1, 2, 3, 15, 29, 30, 31, 59] { + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, minute, 0).unwrap(); + for expr in ["1,2,30 * * * *", "0,1 * * * *", "0,3 9 * * *"] { + assert!( + runs_closer_than(&utc_cron(expr), from, MIN_AGENT_JOB_INTERVAL).is_some(), + "{expr} from :{minute:02}" + ); + } + for expr in ["*/10 * * * *", "0 * * * *", "0,30 9 * * *"] { + assert!( + runs_closer_than(&utc_cron(expr), from, MIN_AGENT_JOB_INTERVAL).is_none(), + "{expr} from :{minute:02}" + ); + } + } +} + +#[test] +fn runs_closer_than_scans_through_the_next_annual_dst_transition() { + let schedule = Schedule::Cron { + expr: "0 0,59 1,3 * * *".into(), + tz: Some("America/New_York".into()), + active_hours: None, + }; + let from = Utc.with_ymd_and_hms(2025, 4, 1, 0, 0, 0).unwrap(); + let hit = runs_closer_than(&schedule, from, MIN_AGENT_JOB_INTERVAL) + .expect("the spring-forward one-minute gap must be found"); + assert_eq!(hit.gap(), ChronoDuration::minutes(1)); + assert_eq!( + hit, + TooFrequent::ConsecutiveRuns { + first: Utc.with_ymd_and_hms(2026, 3, 8, 6, 59, 0).unwrap(), + second: Utc.with_ymd_and_hms(2026, 3, 8, 7, 0, 0).unwrap(), + } + ); +} + +/// The closest pair is reported, not the first one under the floor: `1,3,4,30` +/// has :01 → :03 (2 min) before :03 → :04 (1 min), and the message must name +/// the latter. Ties keep the earliest pair. +#[test] +fn runs_closer_than_reports_the_closest_pair_not_the_first_one_under_the_floor() { + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, 0, 0).unwrap(); + let hit = + runs_closer_than(&utc_cron("1,3,4,30 * * * *"), from, MIN_AGENT_JOB_INTERVAL).unwrap(); + assert_eq!( + hit, + TooFrequent::ConsecutiveRuns { + first: Utc.with_ymd_and_hms(2026, 3, 2, 10, 3, 0).unwrap(), + second: Utc.with_ymd_and_hms(2026, 3, 2, 10, 4, 0).unwrap(), + } + ); + // All gaps equal: the earliest pair wins. + let hit = runs_closer_than(&utc_cron("*/2 * * * *"), from, MIN_AGENT_JOB_INTERVAL).unwrap(); + assert_eq!( + hit, + TooFrequent::ConsecutiveRuns { + first: Utc.with_ymd_and_hms(2026, 3, 2, 10, 2, 0).unwrap(), + second: Utc.with_ymd_and_hms(2026, 3, 2, 10, 4, 0).unwrap(), + } + ); +} + +/// `*/7` is 0,7,…,56: the :56 → :00 step is four minutes, and it counts. +#[test] +fn runs_closer_than_counts_the_wrap_around_gap() { + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, 0, 0).unwrap(); + let hit = runs_closer_than(&utc_cron("*/7 * * * *"), from, MIN_AGENT_JOB_INTERVAL).unwrap(); + assert_eq!(hit.gap(), ChronoDuration::minutes(4)); + assert_eq!( + hit.to_string(), + "fires at 2026-03-02 10:56:00 UTC and again at 2026-03-02 11:00:00 UTC, 4 minutes apart" + ); +} + +#[test] +fn runs_closer_than_is_quiet_at_and_above_the_threshold() { + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, 0, 0).unwrap(); + for expr in [ + "*/5 * * * *", + "*/6 * * * *", + "0 * * * *", + "0 9 * * *", + "0 9 * * 1", + "0 0 1 1 *", + "0 */5 * * * *", + ] { + assert!( + runs_closer_than(&utc_cron(expr), from, MIN_AGENT_JOB_INTERVAL).is_none(), + "{expr}" + ); + } + let every_five = Schedule::Every { every_ms: 300_000 }; + assert!(runs_closer_than(&every_five, from, MIN_AGENT_JOB_INTERVAL).is_none()); + let once = Schedule::At { at: from }; + assert!(runs_closer_than(&once, from, MIN_AGENT_JOB_INTERVAL).is_none()); +} + +#[test] +fn runs_closer_than_reports_a_fixed_interval() { + let from = Utc::now(); + let hit = runs_closer_than( + &Schedule::Every { every_ms: 90_000 }, + from, + MIN_AGENT_JOB_INTERVAL, + ) + .unwrap(); + assert_eq!(hit, TooFrequent::FixedInterval(ChronoDuration::seconds(90))); + assert_eq!(hit.to_string(), "fires every 1 minute 30 seconds"); + let just_under = Schedule::Every { every_ms: 299_999 }; + assert!(runs_closer_than(&just_under, from, MIN_AGENT_JOB_INTERVAL).is_some()); +} + +#[test] +fn runs_closer_than_sees_seconds_level_expressions() { + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, 0, 0).unwrap(); + let hit = runs_closer_than(&utc_cron("*/30 * * * * *"), from, MIN_AGENT_JOB_INTERVAL).unwrap(); + assert_eq!(hit.gap(), ChronoDuration::seconds(30)); +} + +#[test] +fn runs_closer_than_judges_the_effective_cadence_inside_the_active_window() { + let from = Utc.with_ymd_and_hms(2026, 3, 2, 10, 0, 0).unwrap(); + let every_minute_during = |start: &str, end: &str| Schedule::Cron { + expr: "* * * * *".into(), + tz: Some("UTC".into()), + active_hours: Some(ActiveHours { + start: start.into(), + end: end.into(), + }), + }; + // Every minute, but only during one minute of the day: effectively daily. + let daily = every_minute_during("09:00", "09:00"); + assert!(runs_closer_than(&daily, from, MIN_AGENT_JOB_INTERVAL).is_none()); + // A two-minute window lets two adjacent runs through. Skipping ~1,438 + // out-of-window candidates per day exhausts the shared budget long before + // the 1,000-run walk ends; the pair found before that is still reported. + let hit = runs_closer_than( + &every_minute_during("09:00", "09:01"), + from, + MIN_AGENT_JOB_INTERVAL, + ) + .unwrap(); + assert_eq!(hit.gap(), ChronoDuration::minutes(1)); +} + +#[test] +fn runs_closer_than_treats_an_unreadable_expression_as_no_evidence() { + let from = Utc::now(); + assert!(runs_closer_than(&utc_cron("not a cron"), from, MIN_AGENT_JOB_INTERVAL).is_none()); + let bad_tz = Schedule::Cron { + expr: "* * * * *".into(), + tz: Some("Mars/Olympus_Mons".into()), + active_hours: None, + }; + assert!(runs_closer_than(&bad_tz, from, MIN_AGENT_JOB_INTERVAL).is_none()); +} + +#[test] +fn validate_agent_schedule_rejects_tight_schedules_and_names_the_evidence() { + let now = Utc.with_ymd_and_hms(2026, 3, 2, 10, 0, 0).unwrap(); + let err = validate_agent_schedule(&utc_cron("*/3 * * * *"), now) + .unwrap_err() + .to_string(); + assert_eq!( + err, + "Invalid schedule: agent jobs must run at least 5 minutes apart, but this schedule \ + fires at 2026-03-02 10:03:00 UTC and again at 2026-03-02 10:06:00 UTC, 3 minutes apart" + ); + let err = validate_agent_schedule(&Schedule::Every { every_ms: 60_000 }, now) + .unwrap_err() + .to_string(); + assert_eq!( + err, + "Invalid schedule: agent jobs must run at least 5 minutes apart, but this schedule \ + fires every 1 minute" + ); +} + +#[test] +fn validate_agent_schedule_accepts_the_floor_and_keeps_the_generic_checks() { + let now = Utc::now(); + assert!(validate_agent_schedule(&utc_cron("*/5 * * * *"), now).is_ok()); + assert!(validate_agent_schedule(&Schedule::Every { every_ms: 300_000 }, now).is_ok()); + let soon = Schedule::At { + at: now + ChronoDuration::minutes(1), + }; + assert!(validate_agent_schedule(&soon, now).is_ok()); + + // The generic validation runs first: a broken schedule is rejected as + // broken, not as "too frequent". + let err = validate_agent_schedule(&utc_cron("not a cron"), now) + .unwrap_err() + .to_string(); + assert!(err.contains("Invalid cron expression"), "{err}"); + let err = validate_agent_schedule(&Schedule::Every { every_ms: 0 }, now) + .unwrap_err() + .to_string(); + assert!(err.contains("every_ms must be > 0"), "{err}"); + let err = validate_agent_schedule(&Schedule::At { at: now }, now) + .unwrap_err() + .to_string(); + assert!(err.contains("'at' must be in the future"), "{err}"); +} + +#[test] +fn describe_gap_reads_naturally() { + assert_eq!(describe_gap(ChronoDuration::seconds(1)), "1 second"); + assert_eq!(describe_gap(ChronoDuration::seconds(30)), "30 seconds"); + assert_eq!(describe_gap(ChronoDuration::minutes(1)), "1 minute"); + assert_eq!(describe_gap(ChronoDuration::minutes(5)), "5 minutes"); + assert_eq!( + describe_gap(ChronoDuration::seconds(90)), + "1 minute 30 seconds" + ); +} diff --git a/crates/tinyflows-schedule/src/schedule_tests.rs b/crates/tinyflows-schedule/src/schedule_tests.rs new file mode 100644 index 0000000..38f6fbd --- /dev/null +++ b/crates/tinyflows-schedule/src/schedule_tests.rs @@ -0,0 +1,315 @@ +use super::*; +use chrono::TimeZone; + +#[test] +fn next_run_for_schedule_supports_every_and_at() { + let now = Utc::now(); + let every = Schedule::Every { every_ms: 60_000 }; + let next = next_run_for_schedule(&every, now).unwrap(); + assert!(next > now); + + let at = now + ChronoDuration::minutes(10); + let at_schedule = Schedule::At { at }; + let next_at = next_run_for_schedule(&at_schedule, now).unwrap(); + assert_eq!(next_at, at); +} + +#[test] +fn next_run_for_schedule_supports_timezone() { + let from = Utc.with_ymd_and_hms(2026, 2, 16, 0, 0, 0).unwrap(); + let schedule = Schedule::Cron { + expr: "0 9 * * *".into(), + tz: Some("America/Los_Angeles".into()), + active_hours: None, + }; + + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 16, 17, 0, 0).unwrap()); +} + +// ── normalize_expression ──────────────────────────────────────── + +#[test] +fn normalize_expression_accepts_standard_5_field_crontab() { + // 5 fields → seconds column prepended so `cron` crate is happy. + assert_eq!(normalize_expression("0 9 * * *").unwrap(), "0 0 9 * * *"); + assert_eq!( + normalize_expression("*/5 * * * *").unwrap(), + "0 */5 * * * *" + ); +} + +#[test] +fn normalize_expression_accepts_6_and_7_field_crate_native() { + // 6 = second minute hour dom mon dow + assert_eq!(normalize_expression("0 0 9 * * *").unwrap(), "0 0 9 * * *"); + // 7 adds year + assert_eq!( + normalize_expression("0 0 9 * * * 2027").unwrap(), + "0 0 9 * * * 2027" + ); +} + +#[test] +fn normalize_expression_trims_whitespace() { + assert_eq!( + normalize_expression(" 0 9 * * * ").unwrap(), + "0 0 9 * * *" + ); +} + +#[test] +fn normalize_expression_rejects_wrong_field_counts() { + assert!(normalize_expression("").is_err()); + assert!(normalize_expression("* *").is_err()); + assert!(normalize_expression("* * *").is_err()); + assert!(normalize_expression("* * * *").is_err()); + assert!(normalize_expression("* * * * * * * *").is_err()); +} + +// ── next_run_for_schedule ─────────────────────────────────────── + +#[test] +fn next_run_cron_without_tz_uses_local_by_default() { + // Express `from` as local midnight so the expected next-09:00 is always on the + // same calendar day, regardless of the host timezone. A UTC-fixed `from` would + // land at different local times on different machines (e.g. already 10:00 local + // on a UTC+10 host), making the expected date machine-dependent. + let from_local = chrono::Local + .with_ymd_and_hms(2026, 2, 16, 0, 0, 0) + .unwrap(); + let from = from_local.with_timezone(&Utc); + let schedule = Schedule::Cron { + expr: "0 9 * * *".into(), + tz: None, + active_hours: None, + }; + let next = next_run_for_schedule(&schedule, from).unwrap(); + + let expected_local = chrono::Local + .with_ymd_and_hms(2026, 2, 16, 9, 0, 0) + .unwrap(); + assert_eq!(next, expected_local.with_timezone(&Utc)); +} + +#[test] +fn next_run_rejects_invalid_cron_expression() { + let schedule = Schedule::Cron { + expr: "not a cron".into(), + tz: None, + active_hours: None, + }; + let err = next_run_for_schedule(&schedule, Utc::now()).unwrap_err(); + assert!(err.to_string().to_lowercase().contains("invalid")); +} + +#[test] +fn next_run_rejects_invalid_timezone() { + let schedule = Schedule::Cron { + expr: "0 9 * * *".into(), + tz: Some("Not/A_Real_Tz".into()), + active_hours: None, + }; + let err = next_run_for_schedule(&schedule, Utc::now()).unwrap_err(); + assert!( + err.to_string() + .to_lowercase() + .contains("invalid iana timezone") + ); +} + +#[test] +fn next_run_every_zero_is_rejected() { + let schedule = Schedule::Every { every_ms: 0 }; + let err = next_run_for_schedule(&schedule, Utc::now()).unwrap_err(); + assert!(err.to_string().contains("every_ms must be > 0")); +} + +#[test] +fn next_run_at_returns_the_exact_time() { + let at = Utc.with_ymd_and_hms(2026, 3, 1, 12, 0, 0).unwrap(); + let schedule = Schedule::At { at }; + let next = next_run_for_schedule(&schedule, Utc::now()).unwrap(); + assert_eq!(next, at); +} + +// ── validate_schedule ─────────────────────────────────────────── + +#[test] +fn validate_schedule_rejects_past_at_time() { + let now = Utc::now(); + let past = now - ChronoDuration::minutes(5); + let schedule = Schedule::At { at: past }; + let err = validate_schedule(&schedule, now).unwrap_err(); + assert!(err.to_string().contains("'at' must be in the future")); +} + +#[test] +fn validate_schedule_accepts_future_at_time() { + let now = Utc::now(); + let future = now + ChronoDuration::minutes(5); + let schedule = Schedule::At { at: future }; + assert!(validate_schedule(&schedule, now).is_ok()); +} + +#[test] +fn validate_schedule_rejects_every_zero() { + let schedule = Schedule::Every { every_ms: 0 }; + assert!(validate_schedule(&schedule, Utc::now()).is_err()); +} + +#[test] +fn validate_schedule_rejects_every_intervals_that_cannot_advance() { + let now = Utc::now(); + for every_ms in [i64::MAX as u64 + 1, u64::MAX] { + let schedule = Schedule::Every { every_ms }; + assert!(validate_schedule(&schedule, now).is_err()); + } +} + +#[test] +fn validate_schedule_accepts_valid_cron() { + let now = Utc::now(); + let schedule = Schedule::Cron { + expr: "*/5 * * * *".into(), + tz: None, + active_hours: None, + }; + assert!(validate_schedule(&schedule, now).is_ok()); +} + +#[test] +fn validate_schedule_rejects_garbage_cron_expression() { + let schedule = Schedule::Cron { + expr: "not a cron".into(), + tz: None, + active_hours: None, + }; + assert!(validate_schedule(&schedule, Utc::now()).is_err()); +} + +// ── schedule_cron_expression ──────────────────────────────────── + +#[test] +fn schedule_cron_expression_returns_expr_for_cron_variant() { + let s = Schedule::Cron { + expr: "0 9 * * *".into(), + tz: Some("UTC".into()), + active_hours: None, + }; + assert_eq!(schedule_cron_expression(&s).as_deref(), Some("0 9 * * *")); +} + +#[test] +fn schedule_cron_expression_returns_none_for_non_cron_variants() { + assert!(schedule_cron_expression(&Schedule::Every { every_ms: 1000 }).is_none()); + assert!(schedule_cron_expression(&Schedule::At { at: Utc::now() }).is_none()); +} + +#[test] +fn next_run_respects_active_hours() { + // Schedule: every minute + // Active hours: 09:00 - 09:05 + let schedule = Schedule::Cron { + expr: "* * * * *".into(), + tz: Some("UTC".into()), + active_hours: Some(ActiveHours { + start: "09:00".into(), + end: "09:05".into(), + }), + }; + + // If it's 08:00, next run should be 09:00 + let from = Utc.with_ymd_and_hms(2026, 2, 16, 8, 0, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 16, 9, 0, 0).unwrap()); + + // If it's 09:02, next run should be 09:03 + let from = Utc.with_ymd_and_hms(2026, 2, 16, 9, 2, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 16, 9, 3, 0).unwrap()); + + // If it's 09:05, next run should be 09:00 NEXT DAY + let from = Utc.with_ymd_and_hms(2026, 2, 16, 9, 5, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 17, 9, 0, 0).unwrap()); +} + +#[test] +fn next_run_respects_active_hours_spanning_midnight() { + // Active hours: 22:00 - 02:00 + let schedule = Schedule::Cron { + expr: "0 * * * *".into(), // every hour + tz: Some("UTC".into()), + active_hours: Some(ActiveHours { + start: "22:00".into(), + end: "02:00".into(), + }), + }; + + // 20:00 -> 22:00 + let from = Utc.with_ymd_and_hms(2026, 2, 16, 20, 0, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 16, 22, 0, 0).unwrap()); + + // 23:00 -> 00:00 + let from = Utc.with_ymd_and_hms(2026, 2, 16, 23, 0, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 17, 0, 0, 0).unwrap()); + + // 01:00 -> 02:00 + let from = Utc.with_ymd_and_hms(2026, 2, 17, 1, 0, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 17, 2, 0, 0).unwrap()); + + // 03:00 -> 22:00 SAME DAY (since it's early morning) + let from = Utc.with_ymd_and_hms(2026, 2, 17, 3, 0, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 17, 22, 0, 0).unwrap()); +} + +#[test] +fn next_run_respects_active_hours_in_schedule_timezone() { + let schedule = Schedule::Cron { + expr: "0 * * * *".into(), + tz: Some("America/Los_Angeles".into()), + active_hours: Some(ActiveHours { + start: "09:00".into(), + end: "10:00".into(), + }), + }; + + let from = Utc.with_ymd_and_hms(2026, 2, 16, 15, 30, 0).unwrap(); + let next = next_run_for_schedule(&schedule, from).unwrap(); + + assert_eq!(next, Utc.with_ymd_and_hms(2026, 2, 16, 17, 0, 0).unwrap()); +} + +#[test] +fn validate_schedule_rejects_invalid_active_hours() { + let now = Utc::now(); + let schedule = Schedule::Cron { + expr: "* * * * *".into(), + tz: None, + active_hours: Some(ActiveHours { + start: "invalid".into(), + end: "09:00".into(), + }), + }; + assert!(validate_schedule(&schedule, now).is_err()); +} + +#[test] +fn validate_schedule_rejects_invalid_active_hours_end() { + let now = Utc::now(); + let schedule = Schedule::Cron { + expr: "* * * * *".into(), + tz: Some("UTC".into()), + active_hours: Some(ActiveHours { + start: "09:00".into(), + end: "24:00".into(), + }), + }; + let err = validate_schedule(&schedule, now).unwrap_err(); + assert!(err.to_string().contains("active_hours.end")); +} diff --git a/crates/tinyflows-schedule/src/types.rs b/crates/tinyflows-schedule/src/types.rs new file mode 100644 index 0000000..6dd0e79 --- /dev/null +++ b/crates/tinyflows-schedule/src/types.rs @@ -0,0 +1,305 @@ +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Deserializer, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)] +#[serde(rename_all = "lowercase")] +pub enum JobType { + #[default] + Shell, + Agent, + /// A `flows::Flow` schedule trigger binding (issue B2). The job's + /// `command` column carries the bound flow's id (there is no shell + /// command / agent prompt to run); on fire the scheduler publishes + /// `DomainEvent::FlowScheduleTick { flow_id }` instead of running + /// anything itself — `flows::bus::FlowTriggerSubscriber` does the actual + /// dispatch. Created by `flows::ops::flows_set_enabled` (via + /// `cron::add_flow_schedule_job`), never via the `cron_add` agent tool. + Flow, +} + +impl JobType { + pub fn as_str(&self) -> &'static str { + match self { + Self::Shell => "shell", + Self::Agent => "agent", + Self::Flow => "flow", + } + } + + pub fn parse(raw: &str) -> Self { + if raw.eq_ignore_ascii_case("agent") { + Self::Agent + } else if raw.eq_ignore_ascii_case("flow") { + Self::Flow + } else { + Self::Shell + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)] +#[serde(rename_all = "lowercase")] +pub enum SessionTarget { + #[default] + Isolated, + Main, +} + +impl SessionTarget { + pub fn as_str(&self) -> &'static str { + match self { + Self::Isolated => "isolated", + Self::Main => "main", + } + } + + pub fn parse(raw: &str) -> Self { + if raw.eq_ignore_ascii_case("main") { + Self::Main + } else { + Self::Isolated + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct ActiveHours { + pub start: String, + pub end: String, +} + +/// A cron-job schedule. +/// +/// Serializes as an internally-tagged object (`{"kind": "cron", ...}`). +/// Deserializes from **either** that object form **or** a bare cron-expression +/// string like `"0 9 * * 1"` — the bare-string form is treated as +/// `Schedule::Cron { expr, tz: None, active_hours: None }`. +/// +/// The bare-string shorthand exists because agents and some older frontend +/// callers pass `schedule: "0 9 * * 1"` directly instead of the structured +/// object. Accepting it here prevents Sentry issue CORE-RUST-FY +/// ("invalid type: string, expected internally tagged enum Schedule"). +#[derive(Debug, Clone, Serialize, PartialEq, Eq)] +#[serde(tag = "kind", rename_all = "lowercase")] +pub enum Schedule { + Cron { + expr: String, + #[serde(default)] + tz: Option, + #[serde(default)] + active_hours: Option, + }, + At { + at: DateTime, + }, + Every { + every_ms: u64, + }, +} + +impl<'de> Deserialize<'de> for Schedule { + fn deserialize>(deserializer: D) -> Result { + use serde::de::{self, MapAccess, Visitor}; + use std::fmt; + + struct ScheduleVisitor; + + impl<'de> Visitor<'de> for ScheduleVisitor { + type Value = Schedule; + + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + f, + "a cron-schedule object ({{\"kind\":\"cron\",\"expr\":\"...\"}}) \ + or a bare cron-expression string" + ) + } + + /// Accept a bare string as `Schedule::Cron { expr, .. }`. + /// This handles callers that send `schedule: "0 9 * * 1"` directly + /// instead of the structured form. + fn visit_str(self, value: &str) -> Result { + tracing::debug!( + "[cron] Schedule::deserialize: got bare string '{}', \ + coercing to Cron variant", + value + ); + Ok(Schedule::Cron { + expr: value.to_owned(), + tz: None, + active_hours: None, + }) + } + + fn visit_string(self, value: String) -> Result { + tracing::debug!( + "[cron] Schedule::deserialize: got bare string '{}', \ + coercing to Cron variant", + value + ); + Ok(Schedule::Cron { + expr: value, + tz: None, + active_hours: None, + }) + } + + /// Accept the standard internally-tagged object form. + fn visit_map>(self, map: A) -> Result { + // Delegate to the serde-derived tagged-enum logic by + // deserializing from a collected map value. + #[derive(Deserialize)] + #[serde(tag = "kind", rename_all = "lowercase")] + enum ScheduleTagged { + Cron { + expr: String, + #[serde(default)] + tz: Option, + #[serde(default)] + active_hours: Option, + }, + At { + at: DateTime, + }, + Every { + every_ms: u64, + }, + } + + let tagged = + ScheduleTagged::deserialize(de::value::MapAccessDeserializer::new(map))?; + Ok(match tagged { + ScheduleTagged::Cron { + expr, + tz, + active_hours, + } => Schedule::Cron { + expr, + tz, + active_hours, + }, + ScheduleTagged::At { at } => Schedule::At { at }, + ScheduleTagged::Every { every_ms } => Schedule::Every { every_ms }, + }) + } + } + + deserializer.deserialize_any(ScheduleVisitor) + } +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct DeliveryConfig { + #[serde(default)] + pub mode: String, + #[serde(default)] + pub channel: Option, + #[serde(default)] + pub to: Option, + #[serde(default = "default_true")] + pub best_effort: bool, +} + +impl Default for DeliveryConfig { + fn default() -> Self { + Self { + mode: "none".to_string(), + channel: None, + to: None, + best_effort: true, + } + } +} + +fn default_true() -> bool { + true +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CronJob { + pub id: String, + pub expression: String, + pub schedule: Schedule, + pub command: String, + pub prompt: Option, + pub name: Option, + pub job_type: JobType, + pub session_target: SessionTarget, + pub model: Option, + /// Optional built-in agent definition ID (e.g. `"welcome"`, + /// `"morning_briefing"`). When set, the cron scheduler + /// resolves the agent definition from the registry and runs with the + /// definition's prompt, tool allowlist, iteration cap, and model hint + /// instead of the generic `OpenHumanSessionHost::from_config` path. + pub agent_id: Option, + pub enabled: bool, + pub delivery: DeliveryConfig, + pub delete_after_run: bool, + pub created_at: DateTime, + pub next_run: DateTime, + pub last_run: Option>, + pub last_status: Option, + pub last_output: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CronRun { + pub id: i64, + pub job_id: String, + pub started_at: DateTime, + pub finished_at: DateTime, + pub status: String, + pub output: Option, + pub duration_ms: Option, +} + +/// Deserialize a nullable patch field with true double-option semantics: +/// +/// | wire | result | meaning | +/// | --------------- | ------------- | ------------- | +/// | key absent | `None` | no change | +/// | key present `null` | `Some(None)` | clear the value | +/// | key present value | `Some(Some(v))` | set the value | +/// +/// A plain `#[derive(Deserialize)]` on `Option>` collapses the absent +/// and the `null` cases *both* to the outer `None`, so "clear over the wire" +/// (for example, `{"agent_id": null}`) silently deserializes as "no change" — a no-op. Used +/// with `#[serde(default, deserialize_with = "deserialize_double_option")]`, +/// this helper restores the distinction: serde only invokes it when the key is +/// *present*, so a present `null` becomes `Some(None)` and a present value +/// becomes `Some(Some(v))`, while an absent key falls back to the `default` +/// (`None`). +fn deserialize_double_option<'de, D, T>(deserializer: D) -> Result>, D::Error> +where + D: Deserializer<'de>, + T: Deserialize<'de>, +{ + Ok(Some(Option::::deserialize(deserializer)?)) +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct CronJobPatch { + pub schedule: Option, + pub command: Option, + pub prompt: Option, + pub name: Option, + pub enabled: Option, + pub delivery: Option, + pub model: Option, + pub session_target: Option, + pub delete_after_run: Option, + /// `Option>` distinguishes "no change" (`None`) from + /// "clear the agent definition" (`Some(None)`). See + /// [`deserialize_double_option`] for why the custom deserializer is required + /// to honor a wire `null` as a clear rather than a silent no-op. + #[serde( + default, + deserialize_with = "deserialize_double_option", + skip_serializing_if = "Option::is_none" + )] + pub agent_id: Option>, +} + +#[cfg(test)] +#[path = "types_tests.rs"] +mod tests; diff --git a/crates/tinyflows-schedule/src/types_tests.rs b/crates/tinyflows-schedule/src/types_tests.rs new file mode 100644 index 0000000..f7124f5 --- /dev/null +++ b/crates/tinyflows-schedule/src/types_tests.rs @@ -0,0 +1,282 @@ +use super::*; +use chrono::TimeZone; +use serde_json::json; + +// ── JobType ──────────────────────────────────────────────────── + +#[test] +fn job_type_parse_and_as_str_roundtrip() { + assert_eq!(JobType::parse("shell").as_str(), "shell"); + assert_eq!(JobType::parse("agent").as_str(), "agent"); + assert_eq!(JobType::parse("flow").as_str(), "flow"); + // Case-insensitive + assert_eq!(JobType::parse("AGENT"), JobType::Agent); + assert_eq!(JobType::parse("Agent"), JobType::Agent); + assert_eq!(JobType::parse("FLOW"), JobType::Flow); + // Anything unknown falls back to Shell (the default) — guards + // against unexpected legacy DB rows silently turning into Agent. + assert_eq!(JobType::parse(""), JobType::Shell); + assert_eq!(JobType::parse("garbage"), JobType::Shell); +} + +#[test] +fn job_type_default_is_shell() { + assert_eq!(JobType::default(), JobType::Shell); +} + +#[test] +fn job_type_serializes_lowercase() { + assert_eq!(serde_json::to_string(&JobType::Shell).unwrap(), "\"shell\""); + assert_eq!(serde_json::to_string(&JobType::Agent).unwrap(), "\"agent\""); +} + +// ── SessionTarget ────────────────────────────────────────────── + +#[test] +fn session_target_parse_and_as_str_roundtrip() { + assert_eq!(SessionTarget::parse("isolated").as_str(), "isolated"); + assert_eq!(SessionTarget::parse("main").as_str(), "main"); + // Case-insensitive + unknown falls back to Isolated (the default). + assert_eq!(SessionTarget::parse("MAIN"), SessionTarget::Main); + assert_eq!(SessionTarget::parse(""), SessionTarget::Isolated); + assert_eq!(SessionTarget::parse("unknown"), SessionTarget::Isolated); +} + +#[test] +fn session_target_default_is_isolated() { + assert_eq!(SessionTarget::default(), SessionTarget::Isolated); +} + +#[test] +fn session_target_serializes_lowercase() { + assert_eq!( + serde_json::to_string(&SessionTarget::Isolated).unwrap(), + "\"isolated\"" + ); + assert_eq!( + serde_json::to_string(&SessionTarget::Main).unwrap(), + "\"main\"" + ); +} + +#[test] +fn cron_job_patch_preserves_absent_and_explicitly_cleared_agent_ids() { + let unrelated = CronJobPatch { + enabled: Some(false), + ..Default::default() + }; + let value = serde_json::to_value(unrelated).unwrap(); + assert!(value.get("agent_id").is_none()); + + let clear = CronJobPatch { + agent_id: Some(None), + ..Default::default() + }; + let value = serde_json::to_value(clear).unwrap(); + assert_eq!(value["agent_id"], serde_json::Value::Null); +} + +// ── Schedule ─────────────────────────────────────────────────── + +#[test] +fn schedule_cron_variant_roundtrips_with_optional_tz() { + let s = Schedule::Cron { + expr: "0 9 * * *".into(), + tz: Some("America/Los_Angeles".into()), + active_hours: None, + }; + let v = serde_json::to_value(&s).unwrap(); + assert_eq!(v["kind"], "cron"); + assert_eq!(v["expr"], "0 9 * * *"); + assert_eq!(v["tz"], "America/Los_Angeles"); + let back: Schedule = serde_json::from_value(v).unwrap(); + assert_eq!(back, s); +} + +#[test] +fn schedule_cron_variant_accepts_missing_tz() { + let raw = json!({ "kind": "cron", "expr": "*/5 * * * *" }); + let s: Schedule = serde_json::from_value(raw).unwrap(); + assert_eq!( + s, + Schedule::Cron { + expr: "*/5 * * * *".into(), + tz: None, + active_hours: None, + } + ); +} + +#[test] +fn schedule_cron_variant_roundtrips_with_active_hours() { + let s = Schedule::Cron { + expr: "*/15 * * * *".into(), + tz: Some("UTC".into()), + active_hours: Some(ActiveHours { + start: "09:00".into(), + end: "17:30".into(), + }), + }; + let v = serde_json::to_value(&s).unwrap(); + assert_eq!(v["active_hours"]["start"], "09:00"); + assert_eq!(v["active_hours"]["end"], "17:30"); + let back: Schedule = serde_json::from_value(v).unwrap(); + assert_eq!(back, s); +} + +#[test] +fn schedule_at_variant_roundtrips_with_utc_timestamp() { + let at = Utc.with_ymd_and_hms(2027, 1, 15, 12, 0, 0).unwrap(); + let s = Schedule::At { at }; + let v = serde_json::to_value(&s).unwrap(); + assert_eq!(v["kind"], "at"); + let back: Schedule = serde_json::from_value(v).unwrap(); + assert_eq!(back, s); +} + +#[test] +fn schedule_every_variant_roundtrips() { + let s = Schedule::Every { every_ms: 60_000 }; + let v = serde_json::to_value(&s).unwrap(); + assert_eq!(v["kind"], "every"); + assert_eq!(v["every_ms"], 60_000); + let back: Schedule = serde_json::from_value(v).unwrap(); + assert_eq!(back, s); +} + +// ── Schedule bare-string deserialization (CORE-RUST-FY fix) ────── +// Callers (agents, older frontend) sometimes pass a bare cron +// expression string like `"0 9 * * 1"` instead of the structured +// `{"kind":"cron","expr":"0 9 * * 1"}` form. Both must parse. + +#[test] +fn schedule_deserializes_bare_cron_string() { + let s: Schedule = serde_json::from_value(json!("0 9 * * 1")).unwrap(); + assert_eq!( + s, + Schedule::Cron { + expr: "0 9 * * 1".into(), + tz: None, + active_hours: None, + } + ); +} + +#[test] +fn schedule_deserializes_bare_5_field_cron_string() { + let s: Schedule = serde_json::from_str("\"*/5 * * * *\"").unwrap(); + assert_eq!( + s, + Schedule::Cron { + expr: "*/5 * * * *".into(), + tz: None, + active_hours: None, + } + ); +} + +#[test] +fn cron_job_patch_accepts_bare_schedule_string() { + // This is the exact payload shape that triggered CORE-RUST-FY: + // {"schedule": "0 9 * * 1"} + let raw = json!({ "schedule": "0 9 * * 1" }); + let patch: CronJobPatch = serde_json::from_value(raw).unwrap(); + assert_eq!( + patch.schedule, + Some(Schedule::Cron { + expr: "0 9 * * 1".into(), + tz: None, + active_hours: None, + }) + ); +} + +#[test] +fn cron_job_patch_still_accepts_structured_schedule_object() { + let raw = json!({ "schedule": { "kind": "cron", "expr": "0 9 * * 1" } }); + let patch: CronJobPatch = serde_json::from_value(raw).unwrap(); + assert_eq!( + patch.schedule, + Some(Schedule::Cron { + expr: "0 9 * * 1".into(), + tz: None, + active_hours: None, + }) + ); +} + +// ── DeliveryConfig ───────────────────────────────────────────── + +#[test] +fn delivery_config_default_is_none_mode_best_effort() { + let d = DeliveryConfig::default(); + assert_eq!(d.mode, "none"); + assert!(d.channel.is_none()); + assert!(d.to.is_none()); + assert!(d.best_effort, "default best_effort must be true"); +} + +#[test] +fn delivery_config_parses_empty_object_with_defaults() { + // A bare `{}` must deserialize with the `#[serde(default)]` / default + // fn fallbacks — otherwise legacy rows without delivery fields would + // fail to load. + let d: DeliveryConfig = serde_json::from_str("{}").unwrap(); + assert_eq!(d.mode, ""); + assert!(d.channel.is_none()); + assert!(d.to.is_none()); + assert!(d.best_effort, "best_effort must default to true"); +} + +#[test] +fn delivery_config_preserves_best_effort_false_override() { + let raw = json!({ "mode": "channel", "best_effort": false }); + let d: DeliveryConfig = serde_json::from_value(raw).unwrap(); + assert_eq!(d.mode, "channel"); + assert!(!d.best_effort); +} + +// ── CronJobPatch ─────────────────────────────────────────────── + +#[test] +fn cron_job_patch_default_is_all_none() { + let p = CronJobPatch::default(); + assert!(p.schedule.is_none()); + assert!(p.command.is_none()); + assert!(p.prompt.is_none()); + assert!(p.name.is_none()); + assert!(p.enabled.is_none()); + assert!(p.delivery.is_none()); + assert!(p.model.is_none()); + assert!(p.session_target.is_none()); + assert!(p.delete_after_run.is_none()); + assert!(p.agent_id.is_none()); +} + +#[test] +fn patch_agent_id_wire_double_option_semantics() { + // Same fix applied consistently to `agent_id` (its doc + the struct-level + // clearing test already document the Some(None)=clear intent). + let absent: CronJobPatch = serde_json::from_value(json!({})).unwrap(); + assert_eq!(absent.agent_id, None, "absent key means no change"); + let cleared: CronJobPatch = serde_json::from_value(json!({ "agent_id": null })).unwrap(); + assert_eq!( + cleared.agent_id, + Some(None), + "wire null must clear the agent definition" + ); + let set: CronJobPatch = serde_json::from_value(json!({ "agent_id": "welcome" })).unwrap(); + assert_eq!(set.agent_id, Some(Some("welcome".to_string()))); +} + +#[test] +fn cron_job_patch_agent_id_supports_explicit_none_clearing() { + // Option> lets callers distinguish "no change" + // (None) from "clear the agent_id" (Some(None)). + let p = CronJobPatch { + agent_id: Some(None), + ..Default::default() + }; + assert!(p.agent_id.is_some()); + assert!(p.agent_id.as_ref().unwrap().is_none()); +} diff --git a/crates/tinyflows-schedule/src/wire_tests.rs b/crates/tinyflows-schedule/src/wire_tests.rs new file mode 100644 index 0000000..ad02ea6 --- /dev/null +++ b/crates/tinyflows-schedule/src/wire_tests.rs @@ -0,0 +1,104 @@ +//! Literal-JSON fixtures pinning the persisted / RPC wire shape of the +//! schedule types. These strings are what a host's job store and RPC clients +//! already hold; if one stops parsing or re-serialising byte-for-byte, a +//! stored job has been orphaned. + +use crate::{ + ActiveHours, CronJob, CronJobPatch, CronRun, DeliveryConfig, JobType, Schedule, SessionTarget, +}; + +fn assert_bytes(json: &str, expected: &T) +where + T: serde::Serialize + serde::de::DeserializeOwned + PartialEq + std::fmt::Debug, +{ + let parsed: T = serde_json::from_str(json).expect("fixture parses"); + assert_eq!(&parsed, expected); + assert_eq!(serde_json::to_string(&parsed).unwrap(), json); +} + +#[test] +fn cron_schedule_wire_bytes() { + assert_bytes( + r#"{"kind":"cron","expr":"0 9 * * 1","tz":"Europe/London","active_hours":{"start":"09:00","end":"17:30"}}"#, + &Schedule::Cron { + expr: "0 9 * * 1".into(), + tz: Some("Europe/London".into()), + active_hours: Some(ActiveHours { + start: "09:00".into(), + end: "17:30".into(), + }), + }, + ); + assert_bytes( + r#"{"kind":"cron","expr":"*/5 * * * *","tz":null,"active_hours":null}"#, + &Schedule::Cron { + expr: "*/5 * * * *".into(), + tz: None, + active_hours: None, + }, + ); +} + +#[test] +fn at_and_every_schedule_wire_bytes() { + let at: Schedule = + serde_json::from_str(r#"{"kind":"at","at":"2026-02-16T17:00:00Z"}"#).unwrap(); + assert!(matches!(at, Schedule::At { .. })); + assert_eq!( + serde_json::to_string(&at).unwrap(), + r#"{"kind":"at","at":"2026-02-16T17:00:00Z"}"# + ); + assert_bytes( + r#"{"kind":"every","every_ms":60000}"#, + &Schedule::Every { every_ms: 60_000 }, + ); +} + +#[test] +fn legacy_shapes_still_deserialize() { + // Missing optional keys and the bare-string shorthand both read as Cron. + let no_optionals: Schedule = + serde_json::from_str(r#"{"kind":"cron","expr":"0 9 * * *"}"#).unwrap(); + let bare: Schedule = serde_json::from_str(r#""0 9 * * *""#).unwrap(); + let expected = Schedule::Cron { + expr: "0 9 * * *".into(), + tz: None, + active_hours: None, + }; + assert_eq!(no_optionals, expected); + assert_eq!(bare, expected); +} + +#[test] +fn enum_and_delivery_wire_bytes() { + assert_bytes(r#""shell""#, &JobType::Shell); + assert_bytes(r#""agent""#, &JobType::Agent); + assert_bytes(r#""flow""#, &JobType::Flow); + assert_bytes(r#""isolated""#, &SessionTarget::Isolated); + assert_bytes(r#""main""#, &SessionTarget::Main); + assert_bytes( + r#"{"mode":"none","channel":null,"to":null,"best_effort":true}"#, + &DeliveryConfig::default(), + ); +} + +#[test] +fn job_and_run_wire_bytes() { + let job_json = r#"{"id":"j1","expression":"0 9 * * *","schedule":{"kind":"cron","expr":"0 9 * * *","tz":null,"active_hours":null},"command":"echo hi","prompt":null,"name":"morning","job_type":"shell","session_target":"isolated","model":null,"agent_id":null,"enabled":true,"delivery":{"mode":"none","channel":null,"to":null,"best_effort":true},"delete_after_run":false,"created_at":"2026-02-16T00:00:00Z","next_run":"2026-02-16T09:00:00Z","last_run":null,"last_status":null,"last_output":null}"#; + let job: CronJob = serde_json::from_str(job_json).unwrap(); + assert_eq!(serde_json::to_string(&job).unwrap(), job_json); + + let run_json = r#"{"id":7,"job_id":"j1","started_at":"2026-02-16T09:00:00Z","finished_at":"2026-02-16T09:00:02Z","status":"ok","output":"hi","duration_ms":2000}"#; + let run: CronRun = serde_json::from_str(run_json).unwrap(); + assert_eq!(serde_json::to_string(&run).unwrap(), run_json); +} + +#[test] +fn patch_double_option_wire_semantics() { + let absent: CronJobPatch = serde_json::from_str("{}").unwrap(); + assert_eq!(absent.agent_id, None); + let cleared: CronJobPatch = serde_json::from_str(r#"{"agent_id":null}"#).unwrap(); + assert_eq!(cleared.agent_id, Some(None)); + let set: CronJobPatch = serde_json::from_str(r#"{"agent_id":"welcome"}"#).unwrap(); + assert_eq!(set.agent_id, Some(Some("welcome".into()))); +}