car_server_core/
command_scheduler.rs1use std::collections::HashSet;
17use std::time::Duration;
18
19use car_scheduler::os_schedule::LABEL_PREFIX;
20use car_scheduler::{parse_interval, run_command, TaskStore, TaskTrigger};
21
22const POLL_SECS: u64 = 30;
25
26const MAX_EXECUTIONS: usize = 50;
29
30pub fn spawn_command_scheduler() {
33 tokio::spawn(async move {
34 let mut ticker = tokio::time::interval(Duration::from_secs(POLL_SECS));
35 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
36 loop {
37 ticker.tick().await;
38 if let Err(e) = poll_once().await {
39 tracing::warn!(target: "car::scheduler", "[command-scheduler] {e}");
40 }
41 }
42 });
43}
44
45async fn poll_once() -> Result<(), String> {
47 let store = TaskStore::new(&TaskStore::default_path());
48 let tasks = store
50 .try_list()
51 .map_err(|e| format!("could not read task store: {e}"))?;
52 let os_installed: HashSet<String> = car_scheduler::list_installed()
55 .map_err(|e| format!("could not read OS schedules (skipping tick): {e}"))?
56 .into_iter()
57 .collect();
58 let now = chrono::Utc::now();
59
60 for mut task in tasks {
61 if !task.enabled || !task.is_command() || task.trigger != TaskTrigger::Interval {
62 continue;
63 }
64 let label = format!("{LABEL_PREFIX}{}", task.id);
66 if os_installed.contains(&label) {
67 continue;
68 }
69 let interval_secs = parse_interval(&task.schedule);
70 if interval_secs <= 0.0 {
71 continue;
72 }
73 let due = match task.last_run_at {
74 None => true, Some(last) => (now - last).num_seconds() as f64 >= interval_secs,
76 };
77 if !due {
78 continue;
79 }
80 let Some(cmd) = task.command.clone() else {
81 continue;
82 };
83 let exec = run_command(&cmd).await;
84 tracing::info!(
85 target: "car::scheduler",
86 task = %task.id, status = ?exec.status,
87 "[command-scheduler] fired command task"
88 );
89 task.last_run_at = Some(now);
90 task.run_count = task.run_count.saturating_add(1);
91 task.status = exec.status;
92 task.executions.push(exec);
93 if task.executions.len() > MAX_EXECUTIONS {
94 let drop = task.executions.len() - MAX_EXECUTIONS;
95 task.executions.drain(0..drop);
96 }
97 if let Err(e) = store.save(&task) {
98 tracing::warn!(target: "car::scheduler", "[command-scheduler] save {}: {e}", task.id);
99 }
100 }
101 Ok(())
102}