1use anyhow::{Context, Result, anyhow, bail};
2use chrono::{
3 DateTime, Datelike, Duration as ChronoDuration, Local, LocalResult, NaiveDateTime, NaiveTime, TimeZone, Timelike,
4 Utc,
5};
6use humantime::parse_duration as parse_human_duration;
7use regex::Regex;
8use serde::{Deserialize, Serialize};
9use std::collections::{BTreeMap, BTreeSet};
10use std::fmt;
11use std::fs;
12use std::io::Write;
13use std::path::{Path, PathBuf};
14use std::process::Stdio;
15use std::sync::LazyLock;
16use std::sync::atomic::{AtomicU64, Ordering};
17use std::time::Duration;
18use tokio::process::Command;
19
20use crate::notifications::{NotificationEvent, send_global_notification};
21use crate::utils::path::normalize_path;
22use vtcode_commons::VtCodePaths;
23
24pub const MAX_SCHEDULED_TASKS: usize = 50;
25pub const SESSION_TASK_EXPIRY_HOURS: i64 = 72;
26pub const DISABLE_CRON_ENV: &str = "VTCODE_DISABLE_CRON";
27pub const DURABLE_SCHEDULER_RUNTIME_HINT: &str = "Durable tasks fire while VT Code is open, `vtcode schedule serve` is running, or the installed scheduler service is active.";
28
29const SESSION_JITTER_CAP_SECS: u64 = 15 * 60;
30const ONE_SHOT_TOP_OF_HOUR_JITTER_SECS: u64 = 90;
31const CLAIM_STALE_SECS: u64 = 15 * 60;
32const SERVICE_NAME: &str = "vtcode-scheduler";
33const LAUNCHD_LABEL: &str = "com.vtcode.scheduler";
34const DURABLE_STORE_DIR: &str = "scheduled_tasks";
35
36static NEXT_TASK_COUNTER: AtomicU64 = AtomicU64::new(1);
37
38#[cfg(test)]
39mod test_env_overrides {
40 use std::sync::{LazyLock, Mutex};
41
42 static DISABLE_CRON: LazyLock<Mutex<Option<String>>> = LazyLock::new(|| Mutex::new(None));
43
44 pub(super) fn get() -> Option<String> {
45 DISABLE_CRON.lock().unwrap_or_else(|e| e.into_inner()).clone()
46 }
47
48 pub(super) fn set(value: Option<&str>) {
49 if let Ok(mut slot) = DISABLE_CRON.lock() {
50 *slot = value.map(ToString::to_string);
51 }
52 }
53}
54
55#[cfg(test)]
56mod tests;
57
58static REMIND_AT_RE: LazyLock<Regex> = LazyLock::new(|| {
59 Regex::new(r"(?ix)^\s*remind\s+me\s+at\s+(?P<when>.+?)\s+to\s+(?P<prompt>.+)\s*$").expect("remind-at regex")
60});
61
62static REMIND_IN_RE: LazyLock<Regex> = LazyLock::new(|| {
63 Regex::new(
64 r"(?ix)
65 ^\s*in\s+
66 (?P<count>\d+)
67 \s*
68 (?P<unit>minutes|minute|hours|hour|days|day)
69 \s*,?\s*
70 (?P<prompt>.+)
71 \s*$
72 ",
73 )
74 .expect("remind-in regex")
75});
76
77static TIME_ONLY_RE: LazyLock<Regex> = LazyLock::new(|| {
78 Regex::new(r"(?ix)^\s*(?P<hour>\d{1,2})(?::(?P<minute>\d{2}))?\s*(?P<ampm>am|pm)?\s*$").expect("time-only regex")
79});
80
81#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
82#[serde(tag = "kind", rename_all = "snake_case")]
83pub enum ScheduledTaskAction {
84 Prompt { prompt: String },
85 Reminder { message: String },
86}
87
88impl ScheduledTaskAction {
89 #[must_use]
90 pub fn summary(&self) -> &str {
91 match self {
92 Self::Prompt { prompt } => prompt,
93 Self::Reminder { message } => message,
94 }
95 }
96
97 #[must_use]
98 pub fn kind_label(&self) -> &'static str {
99 match self {
100 Self::Prompt { .. } => "prompt",
101 Self::Reminder { .. } => "reminder",
102 }
103 }
104}
105
106#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
107#[serde(tag = "type", rename_all = "snake_case")]
108pub enum ScheduleSpec {
109 Cron5(Cron5),
110 FixedInterval(FixedInterval),
111 OneShot(OneShot),
112}
113
114impl ScheduleSpec {
115 pub fn cron5(expression: impl Into<String>) -> Result<Self> {
116 Ok(Self::Cron5(Cron5::parse(expression)?))
117 }
118
119 pub fn fixed_interval(duration: Duration) -> Result<Self> {
120 Ok(Self::FixedInterval(FixedInterval::from_duration(duration)?))
121 }
122
123 pub fn one_shot(at: DateTime<Utc>) -> Self {
124 Self::OneShot(OneShot { at })
125 }
126
127 #[must_use]
128 pub fn is_recurring(&self) -> bool {
129 !matches!(self, Self::OneShot(_))
130 }
131
132 #[must_use]
133 pub fn human_description(&self) -> String {
134 match self {
135 Self::Cron5(spec) => format!("cron {}", spec.expression),
136 Self::FixedInterval(spec) => spec.human_description(),
137 Self::OneShot(spec) => {
138 format!("once at {}", spec.at.with_timezone(&Local).format("%Y-%m-%d %H:%M"))
139 }
140 }
141 }
142
143 pub fn first_base_fire_at(&self, created_at: DateTime<Utc>) -> Result<Option<DateTime<Utc>>> {
144 let created_local = created_at.with_timezone(&Local);
145 let base_local = match self {
146 Self::Cron5(spec) => spec.next_after(created_local)?,
147 Self::FixedInterval(spec) => Some(created_local + spec.chrono_duration()?),
148 Self::OneShot(spec) => Some(spec.at.with_timezone(&Local)),
149 };
150 Ok(base_local.map(|value| value.with_timezone(&Utc)))
151 }
152
153 pub fn next_base_fire_after(&self, last_base_fire_at: DateTime<Utc>) -> Result<Option<DateTime<Utc>>> {
154 let last_base_local = last_base_fire_at.with_timezone(&Local);
155 let next_local = match self {
156 Self::Cron5(spec) => spec.next_after(last_base_local)?,
157 Self::FixedInterval(spec) => Some(last_base_local + spec.chrono_duration()?),
158 Self::OneShot(_) => None,
159 };
160 Ok(next_local.map(|value| value.with_timezone(&Utc)))
161 }
162
163 fn jittered_fire_at(&self, id: &str, base_fire_at: DateTime<Utc>) -> Result<DateTime<Utc>> {
164 let base_local = base_fire_at.with_timezone(&Local);
165 let hash = stable_hash_u64(id.as_bytes());
166 match self {
167 Self::Cron5(_) | Self::FixedInterval(_) => {
168 let Some(period) = self.nominal_period()? else {
169 return Ok(base_fire_at);
170 };
171 #[allow(
172 clippy::cast_sign_loss,
173 reason = "Intentional compatibility, platform, or test-only suppression."
174 )]
175 let period_secs = period.num_seconds().max(0) as u64;
176 #[allow(
177 clippy::cast_sign_loss,
178 reason = "Intentional compatibility, platform, or test-only suppression."
179 )]
180 let max_delay = (((period_secs as f64) * 0.10).floor()).max(0.0) as u64;
181 let max_delay = max_delay.min(SESSION_JITTER_CAP_SECS);
182 let delay_secs = if max_delay == 0 { 0 } else { hash % (max_delay + 1) };
183 Ok(base_fire_at + ChronoDuration::seconds(delay_secs as i64))
184 }
185 Self::OneShot(_) if matches!(base_local.minute(), 0 | 30) => {
186 let lead_secs = hash % (ONE_SHOT_TOP_OF_HOUR_JITTER_SECS + 1);
187 Ok(base_fire_at - ChronoDuration::seconds(lead_secs as i64))
188 }
189 Self::OneShot(_) => Ok(base_fire_at),
190 }
191 }
192
193 fn nominal_period(&self) -> Result<Option<ChronoDuration>> {
194 match self {
195 Self::FixedInterval(spec) => Ok(Some(spec.chrono_duration()?)),
196 Self::Cron5(spec) => spec.approx_period(),
197 Self::OneShot(_) => Ok(None),
198 }
199 }
200}
201
202#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
203pub struct Cron5 {
204 pub expression: String,
205}
206
207impl Cron5 {
208 pub fn parse(expression: impl Into<String>) -> Result<Self> {
209 let expression = expression.into();
210 ParsedCron::parse(&expression)?;
211 Ok(Self { expression })
212 }
213
214 fn parsed(&self) -> Result<ParsedCron> {
215 ParsedCron::parse(&self.expression)
216 }
217
218 fn next_after(&self, after: DateTime<Local>) -> Result<Option<DateTime<Local>>> {
219 self.parsed()?.next_after(after)
220 }
221
222 fn approx_period(&self) -> Result<Option<ChronoDuration>> {
223 let now = Local::now();
224 let Some(first) = self.next_after(now)? else {
225 return Ok(None);
226 };
227 let Some(second) = self.next_after(first)? else {
228 return Ok(None);
229 };
230 Ok(Some(second - first))
231 }
232}
233
234#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
235pub struct FixedInterval {
236 pub seconds: u64,
237}
238
239impl FixedInterval {
240 pub fn from_duration(duration: Duration) -> Result<Self> {
241 if duration.as_secs() < 60 {
242 bail!("Fixed intervals must be at least 1 minute");
243 }
244 Ok(Self { seconds: duration.as_secs() })
245 }
246
247 pub fn chrono_duration(&self) -> Result<ChronoDuration> {
248 let seconds = i64::try_from(self.seconds).context("interval is too large")?;
249 Ok(ChronoDuration::seconds(seconds))
250 }
251
252 #[must_use]
253 pub fn human_description(&self) -> String {
254 humanize_interval(self.seconds)
255 }
256}
257
258#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
259pub struct OneShot {
260 pub at: DateTime<Utc>,
261}
262
263#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
264pub struct ScheduledTaskDefinition {
265 pub id: String,
266 pub name: String,
267 pub schedule: ScheduleSpec,
268 pub action: ScheduledTaskAction,
269 pub workspace: Option<PathBuf>,
270 pub created_at: DateTime<Utc>,
271 pub expires_at: Option<DateTime<Utc>>,
272}
273
274impl ScheduledTaskDefinition {
275 pub fn new(
276 name: Option<String>,
277 schedule: ScheduleSpec,
278 action: ScheduledTaskAction,
279 workspace: Option<PathBuf>,
280 created_at: DateTime<Utc>,
281 expires_at: Option<DateTime<Utc>>,
282 ) -> Result<Self> {
283 let rendered_name = name
284 .map(|value| value.trim().to_string())
285 .filter(|value| !value.is_empty())
286 .unwrap_or_else(|| summarize_task_name(action.summary()));
287 let id = generate_task_id(&rendered_name, action.summary(), created_at);
288 Ok(Self {
289 id,
290 name: rendered_name,
291 schedule,
292 action,
293 workspace,
294 created_at,
295 expires_at,
296 })
297 }
298}
299
300#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
301#[serde(tag = "state", rename_all = "snake_case")]
302pub enum TaskRunStatus {
303 Triggered,
304 ReminderSent,
305 Success,
306 Failed { message: String },
307}
308
309impl fmt::Display for TaskRunStatus {
310 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
311 match self {
312 Self::Triggered => write!(f, "triggered"),
313 Self::ReminderSent => write!(f, "reminder_sent"),
314 Self::Success => write!(f, "success"),
315 Self::Failed { message } => write!(f, "failed: {message}"),
316 }
317 }
318}
319
320#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
321pub struct ScheduledTaskRuntimeState {
322 pub last_run_at: Option<DateTime<Utc>>,
323 pub next_base_run_at: Option<DateTime<Utc>>,
324 pub next_run_at: Option<DateTime<Utc>>,
325 pub last_status: Option<TaskRunStatus>,
326 pub last_artifact_dir: Option<PathBuf>,
327 pub last_events_file: Option<PathBuf>,
328 pub last_message_file: Option<PathBuf>,
329}
330
331#[derive(Debug, Clone, PartialEq, Eq)]
332pub struct ScheduledTaskRecord {
333 pub definition: ScheduledTaskDefinition,
334 pub runtime: ScheduledTaskRuntimeState,
335}
336
337#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
338pub struct ScheduledTaskSummary {
339 pub id: String,
340 pub name: String,
341 pub action_kind: String,
342 pub schedule: String,
343 pub workspace: Option<PathBuf>,
344 pub recurring: bool,
345 pub next_run_at: Option<DateTime<Utc>>,
346 pub last_run_at: Option<DateTime<Utc>>,
347 pub last_status: Option<String>,
348}
349
350impl ScheduledTaskRecord {
351 fn summary(&self) -> ScheduledTaskSummary {
352 ScheduledTaskSummary {
353 id: self.definition.id.clone(),
354 name: self.definition.name.clone(),
355 action_kind: self.definition.action.kind_label().to_string(),
356 schedule: self.definition.schedule.human_description(),
357 workspace: self.definition.workspace.clone(),
358 recurring: self.definition.schedule.is_recurring(),
359 next_run_at: self.runtime.next_run_at,
360 last_run_at: self.runtime.last_run_at,
361 last_status: self.runtime.last_status.as_ref().map(ToString::to_string),
362 }
363 }
364}
365
366#[derive(Debug, Clone, PartialEq, Eq)]
367pub struct DueSessionPrompt {
368 pub id: String,
369 pub name: String,
370 pub prompt: String,
371}
372
373#[derive(Debug, Default, Clone)]
374pub struct SessionScheduler {
375 tasks: BTreeMap<String, ScheduledTaskRecord>,
376}
377
378impl SessionScheduler {
379 #[must_use]
380 pub fn new() -> Self {
381 Self { tasks: BTreeMap::new() }
382 }
383
384 #[must_use]
385 pub fn len(&self) -> usize {
386 self.tasks.len()
387 }
388
389 #[must_use]
390 pub fn is_empty(&self) -> bool {
391 self.tasks.is_empty()
392 }
393
394 pub fn create_prompt_task(
395 &mut self,
396 name: Option<String>,
397 prompt: String,
398 schedule: ScheduleSpec,
399 created_at: DateTime<Utc>,
400 ) -> Result<ScheduledTaskSummary> {
401 self.ensure_capacity()?;
402 let expires_at = schedule
403 .is_recurring()
404 .then(|| created_at + ChronoDuration::hours(SESSION_TASK_EXPIRY_HOURS));
405 let definition = ScheduledTaskDefinition::new(
406 name,
407 schedule,
408 ScheduledTaskAction::Prompt { prompt },
409 None,
410 created_at,
411 expires_at,
412 )?;
413 let runtime = initialize_runtime_state(&definition)?;
414 let record = ScheduledTaskRecord { definition: definition.clone(), runtime };
415 let summary = record.summary();
416 self.tasks.insert(definition.id.clone(), record);
417 Ok(summary)
418 }
419
420 pub fn list(&self) -> Vec<ScheduledTaskSummary> {
421 self.tasks.values().map(ScheduledTaskRecord::summary).collect()
422 }
423
424 pub fn delete(&mut self, query: &str) -> Option<ScheduledTaskSummary> {
425 let query = query.trim();
426 if query.is_empty() {
427 return None;
428 }
429 if let Some(record) = self.tasks.remove(query) {
430 return Some(record.summary());
431 }
432 let key = self
433 .tasks
434 .iter()
435 .find_map(|(id, record)| record.definition.name.eq_ignore_ascii_case(query).then(|| id.clone()))?;
436 self.tasks.remove(&key).map(|record| record.summary())
437 }
438
439 pub fn collect_due_prompts(&mut self, now: DateTime<Utc>) -> Result<Vec<DueSessionPrompt>> {
440 let mut due = Vec::new();
441 let mut completed = Vec::new();
442 for record in self.tasks.values_mut() {
443 let Some(next_run_at) = record.runtime.next_run_at else {
444 continue;
445 };
446 if now < next_run_at {
447 continue;
448 }
449
450 if let Some(prompt) = due_session_prompt(record) {
451 due.push(prompt);
452 }
453
454 if advance_record_runtime(record, now, TaskRunStatus::Triggered)? {
455 completed.push(record.definition.id.clone());
456 }
457 }
458
459 for id in completed {
460 self.tasks.remove(&id);
461 }
462
463 Ok(due)
464 }
465
466 fn ensure_capacity(&self) -> Result<()> {
467 if self.tasks.len() >= MAX_SCHEDULED_TASKS {
468 bail!("A session can hold at most {MAX_SCHEDULED_TASKS} scheduled tasks");
469 }
470 Ok(())
471 }
472}
473
474#[derive(Debug, Clone, PartialEq, Eq)]
475pub enum SessionLanguageCommand {
476 CreateOneShotPrompt { prompt: String, run_at: DateTime<Utc> },
477 ListTasks,
478 CancelTask { query: String },
479}
480
481#[derive(Debug, Clone, PartialEq, Eq)]
482pub struct ScheduleCreateInput {
483 pub name: Option<String>,
484 pub prompt: Option<String>,
485 pub reminder: Option<String>,
486 pub every: Option<String>,
487 pub cron: Option<String>,
488 pub at: Option<String>,
489 pub workspace: Option<PathBuf>,
490}
491
492impl ScheduleCreateInput {
493 pub fn build_definition(
494 self,
495 now: DateTime<Local>,
496 default_workspace: Option<PathBuf>,
497 ) -> Result<ScheduledTaskDefinition> {
498 let action = match (self.prompt, self.reminder) {
499 (Some(prompt), None) => ScheduledTaskAction::Prompt { prompt },
500 (None, Some(reminder)) => ScheduledTaskAction::Reminder { message: reminder },
501 (Some(_), Some(_)) => bail!("Choose either --prompt or --reminder"),
502 (None, None) => bail!("One of --prompt or --reminder is required"),
503 };
504
505 let schedule = match (self.every, self.cron, self.at) {
506 (Some(raw), None, None) => {
507 let duration = parse_human_duration(raw.trim())
508 .with_context(|| format!("Invalid --every duration: {}", raw.trim()))?;
509 ScheduleSpec::fixed_interval(duration)?
510 }
511 (None, Some(expression), None) => ScheduleSpec::cron5(expression)?,
512 (None, None, Some(raw)) => ScheduleSpec::one_shot(parse_local_datetime(raw.trim(), now)?),
513 _ => bail!("Choose exactly one of --every, --cron, or --at"),
514 };
515
516 let workspace = match (&action, self.workspace.or(default_workspace)) {
517 (ScheduledTaskAction::Prompt { .. }, Some(path)) => Some(path),
518 (ScheduledTaskAction::Prompt { .. }, None) => {
519 bail!("Prompt tasks require a workspace (pass --workspace or create from chat)")
520 }
521 (ScheduledTaskAction::Reminder { .. }, path) => path,
522 }
523 .map(|path| resolve_scheduled_workspace_path(&path))
524 .transpose()?;
525
526 if matches!(action, ScheduledTaskAction::Prompt { .. }) && workspace.as_ref().is_some_and(|path| !path.is_dir())
527 {
528 let workspace = workspace.as_ref().context("prompt workspace should exist")?;
529 bail!("Prompt task workspace does not exist or is not a directory: {}", workspace.display());
530 }
531
532 ScheduledTaskDefinition::new(self.name, schedule, action, workspace, now.with_timezone(&Utc), None)
533 }
534}
535
536fn resolve_scheduled_workspace_path(path: &Path) -> Result<PathBuf> {
537 resolve_scheduled_workspace_path_with_home(path, dirs::home_dir().as_deref())
538}
539
540fn resolve_scheduled_workspace_path_with_home(path: &Path, home_dir: Option<&Path>) -> Result<PathBuf> {
541 let expanded = expand_scheduled_workspace_home(path, home_dir);
542 let absolute = if expanded.is_absolute() {
543 expanded
544 } else {
545 std::env::current_dir()
546 .context("Failed to resolve current directory for scheduled task workspace")?
547 .join(expanded)
548 };
549 Ok(normalize_path(&absolute))
550}
551
552fn expand_scheduled_workspace_home(path: &Path, home_dir: Option<&Path>) -> PathBuf {
553 let Some(raw) = path.to_str() else {
554 return path.to_path_buf();
555 };
556
557 if raw == "~" {
558 return home_dir.map(Path::to_path_buf).unwrap_or_else(|| path.to_path_buf());
559 }
560
561 if let Some(rest) = raw.strip_prefix("~/")
562 && let Some(home_dir) = home_dir
563 {
564 return home_dir.join(rest);
565 }
566
567 path.to_path_buf()
568}
569
570#[derive(Debug, Clone)]
571pub struct SchedulerPaths {
572 pub config_root: PathBuf,
573 pub data_root: PathBuf,
574}
575
576impl SchedulerPaths {
577 pub fn new_default() -> Result<Self> {
578 let paths = VtCodePaths::resolve().context("Failed to resolve VT Code paths")?;
579 let config_root = paths.config_path(DURABLE_STORE_DIR)?;
580 let data_root = paths.state_path(DURABLE_STORE_DIR)?;
581 Ok(Self { config_root, data_root })
582 }
583
584 #[must_use]
585 pub fn tasks_dir(&self) -> PathBuf {
586 self.config_root.join("tasks")
587 }
588
589 #[must_use]
590 pub fn state_dir(&self) -> PathBuf {
591 self.data_root.join("state")
592 }
593
594 #[must_use]
595 pub fn claims_dir(&self) -> PathBuf {
596 self.data_root.join("claims")
597 }
598
599 #[must_use]
600 pub fn runs_dir(&self) -> PathBuf {
601 self.data_root.join("runs")
602 }
603
604 pub fn ensure_dirs(&self) -> Result<()> {
605 for dir in [
606 &self.config_root,
607 &self.data_root,
608 &self.tasks_dir(),
609 &self.state_dir(),
610 &self.claims_dir(),
611 &self.runs_dir(),
612 ] {
613 VtCodePaths::ensure_user_dir(dir)
614 .with_context(|| format!("Failed to create scheduler directory {}", dir.display()))?;
615 }
616 Ok(())
617 }
618
619 #[must_use]
620 pub fn definition_path(&self, id: &str) -> PathBuf {
621 self.tasks_dir().join(format!("{id}.toml"))
622 }
623
624 #[must_use]
625 pub fn runtime_path(&self, id: &str) -> PathBuf {
626 self.state_dir().join(format!("{id}.json"))
627 }
628
629 #[must_use]
630 pub fn claim_path(&self, id: &str) -> PathBuf {
631 self.claims_dir().join(format!("{id}.claim"))
632 }
633
634 #[must_use]
635 pub fn artifact_dir(&self, id: &str, run_at: DateTime<Utc>) -> PathBuf {
636 self.runs_dir().join(id).join(run_at.format("%Y%m%dT%H%M%SZ").to_string())
637 }
638}
639
640#[derive(Debug, Clone)]
641pub struct DurableTaskStore {
642 paths: SchedulerPaths,
643}
644
645impl DurableTaskStore {
646 #[must_use]
647 pub fn with_paths(paths: SchedulerPaths) -> Self {
648 Self { paths }
649 }
650
651 pub fn new_default() -> Result<Self> {
652 Ok(Self { paths: SchedulerPaths::new_default()? })
653 }
654
655 #[must_use]
656 pub fn paths(&self) -> &SchedulerPaths {
657 &self.paths
658 }
659
660 pub fn create(&self, definition: ScheduledTaskDefinition) -> Result<ScheduledTaskSummary> {
661 self.paths.ensure_dirs()?;
662 let current_count = self.definition_paths()?.len();
663 if current_count >= MAX_SCHEDULED_TASKS {
664 bail!("VT Code supports at most {MAX_SCHEDULED_TASKS} durable scheduled tasks");
665 }
666
667 let runtime = initialize_runtime_state(&definition)?;
668 self.write_definition(&definition)?;
669 self.write_runtime(&definition.id, &runtime)?;
670 Ok(ScheduledTaskRecord { definition, runtime }.summary())
671 }
672
673 pub fn create_from_input(
674 &self,
675 input: ScheduleCreateInput,
676 now: DateTime<Local>,
677 default_workspace: Option<PathBuf>,
678 ) -> Result<ScheduledTaskSummary> {
679 let definition = input.build_definition(now, default_workspace)?;
680 self.create(definition)
681 }
682
683 pub fn list(&self) -> Result<Vec<ScheduledTaskSummary>> {
684 let mut records = self.load_records()?;
685 records.sort_by_key(|record| record.runtime.next_run_at);
686 Ok(records.into_iter().map(|record| record.summary()).collect())
687 }
688
689 pub fn delete(&self, id: &str) -> Result<Option<ScheduledTaskSummary>> {
690 let Some(record) = self.load_record(id)? else {
691 return Ok(None);
692 };
693 let _ = fs::remove_file(self.paths.definition_path(id));
694 let _ = fs::remove_file(self.paths.runtime_path(id));
695 let _ = fs::remove_file(self.paths.claim_path(id));
696 Ok(Some(record.summary()))
697 }
698
699 pub fn load_record(&self, id: &str) -> Result<Option<ScheduledTaskRecord>> {
700 let definition_path = self.paths.definition_path(id);
701 if !definition_path.exists() {
702 return Ok(None);
703 }
704 let definition = read_definition(&definition_path)?;
705 let runtime = match self.read_runtime(&definition.id)? {
706 Some(runtime) => runtime,
707 None => initialize_runtime_state(&definition)?,
708 };
709 Ok(Some(ScheduledTaskRecord { definition, runtime }))
710 }
711
712 pub fn update_runtime(&self, record: &ScheduledTaskRecord) -> Result<()> {
713 self.write_runtime(&record.definition.id, &record.runtime)
714 }
715
716 async fn load_records_async(&self) -> Result<Vec<ScheduledTaskRecord>> {
717 let store = self.clone();
718 tokio::task::spawn_blocking(move || store.load_records())
719 .await
720 .context("Scheduled task loading task panicked")?
721 }
722
723 async fn update_runtime_async(&self, record: ScheduledTaskRecord) -> Result<()> {
724 let store = self.clone();
725 tokio::task::spawn_blocking(move || store.update_runtime(&record))
726 .await
727 .context("Scheduled task runtime update task panicked")?
728 }
729
730 fn load_records(&self) -> Result<Vec<ScheduledTaskRecord>> {
731 self.paths.ensure_dirs()?;
732 let mut records = Vec::new();
733 for definition_path in self.definition_paths()? {
734 let definition = read_definition(&definition_path)?;
735 let runtime = self
736 .read_runtime(&definition.id)?
737 .unwrap_or(initialize_runtime_state(&definition)?);
738 records.push(ScheduledTaskRecord { definition, runtime });
739 }
740 Ok(records)
741 }
742
743 fn definition_paths(&self) -> Result<Vec<PathBuf>> {
744 self.paths.ensure_dirs()?;
745 let mut paths = Vec::new();
746 for entry in fs::read_dir(self.paths.tasks_dir())
747 .with_context(|| format!("Failed to read {}", self.paths.tasks_dir().display()))?
748 {
749 let entry = entry?;
750 let path = entry.path();
751 if path.extension().and_then(|value| value.to_str()) == Some("toml") {
752 paths.push(path);
753 }
754 }
755 paths.sort();
756 Ok(paths)
757 }
758
759 fn write_definition(&self, definition: &ScheduledTaskDefinition) -> Result<()> {
760 let serialized = toml::to_string_pretty(definition).context("Failed to serialize task definition")?;
761 atomic_write(&self.paths.definition_path(&definition.id), serialized.as_bytes())
762 }
763
764 fn read_runtime(&self, id: &str) -> Result<Option<ScheduledTaskRuntimeState>> {
765 let path = self.paths.runtime_path(id);
766 if !path.exists() {
767 return Ok(None);
768 }
769 let raw = fs::read_to_string(&path).with_context(|| format!("Failed to read {}", path.display()))?;
770 let runtime = serde_json::from_str(&raw).with_context(|| format!("Failed to parse {}", path.display()))?;
771 Ok(Some(runtime))
772 }
773
774 fn write_runtime(&self, id: &str, runtime: &ScheduledTaskRuntimeState) -> Result<()> {
775 let serialized = serde_json::to_vec_pretty(runtime).context("Failed to serialize runtime state")?;
776 atomic_write(&self.paths.runtime_path(id), &serialized)
777 }
778}
779
780#[derive(Debug, Clone)]
781pub struct SchedulerDaemon {
782 store: DurableTaskStore,
783 executable_path: PathBuf,
784}
785
786impl SchedulerDaemon {
787 #[must_use]
788 pub fn new(store: DurableTaskStore, executable_path: PathBuf) -> Self {
789 Self { store, executable_path }
790 }
791
792 pub async fn serve_forever(&self) -> Result<()> {
793 loop {
794 self.run_due_tasks_once().await?;
795 tokio::time::sleep(Duration::from_secs(1)).await;
796 }
797 }
798
799 pub async fn run_due_tasks_once(&self) -> Result<usize> {
800 let now = Utc::now();
801 let mut records = self.store.load_records_async().await?;
802 let mut executed = 0usize;
803
804 records.sort_by_key(|record| record.runtime.next_run_at);
805 for mut record in records {
806 let Some(next_run_at) = record.runtime.next_run_at else {
807 continue;
808 };
809 if now < next_run_at {
810 continue;
811 }
812 if !try_acquire_claim_async(self.store.paths().clone(), record.definition.id.clone()).await? {
813 continue;
814 }
815
816 let result = self.execute_record(&record, now).await;
817 let release_result = release_claim_async(self.store.paths().clone(), record.definition.id.clone()).await;
818 let run_outcome = result?;
819 release_result?;
820
821 apply_run_outcome(&mut record, run_outcome)?;
822 self.store.update_runtime_async(record).await?;
823 executed = executed.saturating_add(1);
824 }
825
826 Ok(executed)
827 }
828
829 async fn execute_record(&self, record: &ScheduledTaskRecord, run_at: DateTime<Utc>) -> Result<RunOutcome> {
830 match &record.definition.action {
831 ScheduledTaskAction::Reminder { message } => {
832 let notification = send_global_notification(NotificationEvent::IdlePrompt {
833 title: format!("Scheduled reminder: {}", record.definition.name),
834 message: message.clone(),
835 })
836 .await;
837 let status = match notification {
838 Ok(()) => TaskRunStatus::ReminderSent,
839 Err(error) => TaskRunStatus::Failed {
840 message: format!("failed to send reminder notification: {error:#}"),
841 },
842 };
843 Ok(RunOutcome {
844 ran_at: run_at,
845 status,
846 artifact_dir: None,
847 events_file: None,
848 last_message_file: None,
849 })
850 }
851 ScheduledTaskAction::Prompt { prompt } => {
852 let artifact_dir = self.store.paths().artifact_dir(&record.definition.id, run_at);
853 let events_file = artifact_dir.join("events.jsonl");
854 let last_message_file = artifact_dir.join("last-message.txt");
855 let workspace_label = scheduled_workspace_label(&record.definition);
856 let workspace = resolve_scheduled_task_workspace(&record.definition);
857 let execution = async {
858 VtCodePaths::ensure_user_dir(&artifact_dir)
859 .with_context(|| format!("Failed to create run artifact dir {}", artifact_dir.display()))?;
860
861 let mut command = Command::new(&self.executable_path);
862 command
863 .arg("exec")
864 .arg("--events")
865 .arg(&events_file)
866 .arg("--last-message-file")
867 .arg(&last_message_file)
868 .arg(prompt)
869 .stdin(Stdio::null())
870 .stdout(Stdio::null())
871 .stderr(Stdio::null());
872
873 if let Some(workspace) = workspace? {
874 command.current_dir(&workspace);
875 }
876
877 command.status().await.with_context(|| {
878 format!(
879 "Failed to spawn scheduled VT Code exec for task {} using {} in {}",
880 record.definition.id,
881 self.executable_path.display(),
882 workspace_label
883 )
884 })
885 }
886 .await;
887 let run_status = match execution {
888 Ok(status) if status.success() => TaskRunStatus::Success,
889 Ok(status) => TaskRunStatus::Failed {
890 message: format!("vtcode exec exited with {status}"),
891 },
892 Err(error) => TaskRunStatus::Failed { message: format!("{error:#}") },
893 };
894 Ok(RunOutcome {
895 ran_at: run_at,
896 status: run_status,
897 artifact_dir: artifact_dir.is_dir().then_some(artifact_dir),
898 events_file: events_file.is_file().then_some(events_file),
899 last_message_file: last_message_file.is_file().then_some(last_message_file),
900 })
901 }
902 }
903 }
904}
905
906#[derive(Debug, Clone, Copy, PartialEq, Eq)]
907pub enum ServiceManager {
908 Launchd,
909 SystemdUser,
910}
911
912impl ServiceManager {
913 #[must_use]
914 pub fn current() -> Option<Self> {
915 #[cfg(target_os = "macos")]
916 {
917 return Some(Self::Launchd);
918 }
919 #[cfg(all(unix, not(target_os = "macos")))]
920 {
921 return Some(Self::SystemdUser);
922 }
923 #[cfg_attr(
924 unix,
925 allow(
926 unreachable_code,
927 reason = "Platform-specific returns make the fallback unreachable on Unix."
928 )
929 )]
930 None
931 }
932}
933
934#[derive(Debug, Clone)]
935pub struct ServiceInstallPlan {
936 pub manager: ServiceManager,
937 pub path: PathBuf,
938 pub contents: String,
939}
940
941pub fn render_service_install_plan(executable_path: &Path) -> Result<ServiceInstallPlan> {
942 let manager = ServiceManager::current()
943 .ok_or_else(|| anyhow!("Durable scheduler services are unsupported on this platform"))?;
944 let path = match manager {
945 ServiceManager::Launchd => dirs::home_dir()
946 .ok_or_else(|| anyhow!("Failed to resolve home directory"))?
947 .join("Library/LaunchAgents")
948 .join(format!("{LAUNCHD_LABEL}.plist")),
949 ServiceManager::SystemdUser => dirs::home_dir()
950 .ok_or_else(|| anyhow!("Failed to resolve home directory"))?
951 .join(".config/systemd/user")
952 .join(format!("{SERVICE_NAME}.service")),
953 };
954 let contents = match manager {
955 ServiceManager::Launchd => render_launchd_plist(executable_path),
956 ServiceManager::SystemdUser => render_systemd_unit(executable_path),
957 };
958 Ok(ServiceInstallPlan { manager, path, contents })
959}
960
961pub fn install_service_file(executable_path: &Path) -> Result<ServiceInstallPlan> {
962 let plan = render_service_install_plan(executable_path)?;
963 if let Some(parent) = plan.path.parent() {
964 VtCodePaths::ensure_user_dir(parent).with_context(|| format!("Failed to create {}", parent.display()))?;
965 }
966 atomic_write(&plan.path, plan.contents.as_bytes())?;
967 Ok(plan)
968}
969
970pub fn uninstall_service_file() -> Result<Option<(ServiceManager, PathBuf, bool)>> {
971 let Some(manager) = ServiceManager::current() else {
972 return Ok(None);
973 };
974 let path = render_service_install_plan(Path::new("/tmp/vtcode"))?.path;
975 if path.exists() {
976 fs::remove_file(&path).with_context(|| format!("Failed to remove {}", path.display()))?;
977 return Ok(Some((manager, path, true)));
978 }
979 Ok(Some((manager, path, false)))
980}
981
982#[must_use]
983pub fn render_launchd_plist(executable_path: &Path) -> String {
984 format!(
985 r#"<?xml version="1.0" encoding="UTF-8"?>
986<!DOCTYPE plist PUBLIC "-//Apple//DTD PLIST 1.0//EN" "http://www.apple.com/DTDs/PropertyList-1.0.dtd">
987<plist version="1.0">
988 <dict>
989 <key>Label</key>
990 <string>{LAUNCHD_LABEL}</string>
991 <key>ProgramArguments</key>
992 <array>
993 <string>{}</string>
994 <string>schedule</string>
995 <string>serve</string>
996 </array>
997 <key>RunAtLoad</key>
998 <true/>
999 <key>KeepAlive</key>
1000 <true/>
1001 </dict>
1002</plist>
1003"#,
1004 xml_escape(executable_path.display().to_string())
1005 )
1006}
1007
1008#[must_use]
1009pub fn render_systemd_unit(executable_path: &Path) -> String {
1010 format!(
1011 "[Unit]\nDescription=VT Code scheduler\n\n[Service]\nType=simple\nExecStart={} schedule serve\nRestart=always\nRestartSec=5\n\n[Install]\nWantedBy=default.target\n",
1012 shell_words::quote(executable_path.to_string_lossy().as_ref())
1013 )
1014}
1015
1016pub fn scheduled_tasks_enabled(enabled_in_config: bool) -> bool {
1017 #[cfg(test)]
1018 if let Some(value) = test_env_overrides::get() {
1019 let normalized = value.trim().to_ascii_lowercase();
1020 if matches!(normalized.as_str(), "1" | "true" | "yes" | "on") {
1021 return false;
1022 }
1023 }
1024
1025 if let Ok(value) = std::env::var(DISABLE_CRON_ENV) {
1026 let normalized = value.trim().to_ascii_lowercase();
1027 if matches!(normalized.as_str(), "1" | "true" | "yes" | "on") {
1028 return false;
1029 }
1030 }
1031 enabled_in_config
1032}
1033
1034#[must_use]
1035pub fn durable_task_is_overdue(
1036 next_run_at: Option<DateTime<Utc>>,
1037 last_run_at: Option<DateTime<Utc>>,
1038 has_last_status: bool,
1039 now: DateTime<Utc>,
1040) -> bool {
1041 !has_last_status && last_run_at.is_none() && next_run_at.is_some_and(|next_run| next_run <= now)
1042}
1043
1044pub fn parse_session_language_command(input: &str, now: DateTime<Local>) -> Option<Result<SessionLanguageCommand>> {
1045 let trimmed = input.trim();
1046 if trimmed.is_empty() {
1047 return None;
1048 }
1049
1050 let normalized = trimmed.trim_end_matches(['?', '.', '!']).to_ascii_lowercase();
1051 if normalized == "what scheduled tasks do i have" {
1052 return Some(Ok(SessionLanguageCommand::ListTasks));
1053 }
1054
1055 if let Some(captures) = REMIND_AT_RE.captures(trimmed) {
1056 let when = captures.name("when").map(|value| value.as_str()).unwrap_or_default();
1057 let prompt = captures
1058 .name("prompt")
1059 .map(|value| value.as_str().trim().to_string())
1060 .unwrap_or_default();
1061 return Some(
1062 parse_local_datetime(when, now)
1063 .map(|run_at| SessionLanguageCommand::CreateOneShotPrompt { prompt, run_at }),
1064 );
1065 }
1066
1067 if let Some(captures) = REMIND_IN_RE.captures(trimmed) {
1068 let count = captures["count"].parse::<i64>().ok()?;
1069 let prompt = captures["prompt"].trim().to_string();
1070 let delta = match captures["unit"].to_ascii_lowercase().as_str() {
1071 "minute" | "minutes" => ChronoDuration::minutes(count),
1072 "hour" | "hours" => ChronoDuration::hours(count),
1073 "day" | "days" => ChronoDuration::days(count),
1074 _ => return None,
1075 };
1076 return Some(Ok(SessionLanguageCommand::CreateOneShotPrompt {
1077 prompt,
1078 run_at: (now + delta).with_timezone(&Utc),
1079 }));
1080 }
1081
1082 if let Some(query) = trimmed.strip_prefix("cancel ") {
1083 return Some(Ok(SessionLanguageCommand::CancelTask { query: query.trim().to_string() }));
1084 }
1085
1086 None
1087}
1088
1089pub fn parse_schedule_create_args(args: &str) -> Result<ScheduleCreateInput> {
1090 let tokens = shell_words::split(args).with_context(|| format!("Failed to parse arguments: {args}"))?;
1091 parse_schedule_create_tokens(&tokens)
1092}
1093
1094pub fn parse_schedule_create_tokens(tokens: &[String]) -> Result<ScheduleCreateInput> {
1095 let mut name = None;
1096 let mut prompt = None;
1097 let mut reminder = None;
1098 let mut every = None;
1099 let mut cron = None;
1100 let mut at = None;
1101 let mut workspace = None;
1102 let mut index = 0usize;
1103
1104 while index < tokens.len() {
1105 let token = &tokens[index];
1106 let (flag, inline_value) = if let Some((left, right)) = token.split_once('=') {
1107 (left, Some(right.to_string()))
1108 } else {
1109 (token.as_str(), None)
1110 };
1111
1112 let take_value = |idx: &mut usize| -> Result<String> {
1113 if let Some(value) = inline_value.clone() {
1114 return Ok(value);
1115 }
1116 let Some(value) = tokens.get(*idx + 1) else {
1117 bail!("Missing value for {flag}");
1118 };
1119 *idx += 1;
1120 Ok(value.clone())
1121 };
1122
1123 match flag {
1124 "--name" => name = Some(take_value(&mut index)?),
1125 "--prompt" => prompt = Some(take_value(&mut index)?),
1126 "--reminder" => reminder = Some(take_value(&mut index)?),
1127 "--every" => every = Some(take_value(&mut index)?),
1128 "--cron" => cron = Some(take_value(&mut index)?),
1129 "--at" => at = Some(take_value(&mut index)?),
1130 "--workspace" => workspace = Some(PathBuf::from(take_value(&mut index)?)),
1131 "--help" | "help" => {
1132 bail!(
1133 "Usage: /schedule create --prompt <text>|--reminder <text> --every <dur>|--cron <expr>|--at <time> [--name <label>] [--workspace <path>]"
1134 );
1135 }
1136 _ => bail!("Unknown option: {token}"),
1137 }
1138 index += 1;
1139 }
1140
1141 Ok(ScheduleCreateInput { name, prompt, reminder, every, cron, at, workspace })
1142}
1143
1144fn initialize_runtime_state(definition: &ScheduledTaskDefinition) -> Result<ScheduledTaskRuntimeState> {
1145 let next_base_run_at = definition.schedule.first_base_fire_at(definition.created_at)?;
1146 let next_run_at = next_base_run_at
1147 .map(|base| definition.schedule.jittered_fire_at(&definition.id, base))
1148 .transpose()?;
1149 Ok(ScheduledTaskRuntimeState {
1150 next_base_run_at,
1151 next_run_at,
1152 ..ScheduledTaskRuntimeState::default()
1153 })
1154}
1155
1156#[derive(Debug, Clone, Copy)]
1157struct NextScheduledRun {
1158 base_fire_at: DateTime<Utc>,
1159 fire_at: DateTime<Utc>,
1160}
1161
1162fn due_session_prompt(record: &ScheduledTaskRecord) -> Option<DueSessionPrompt> {
1163 let ScheduledTaskAction::Prompt { prompt } = &record.definition.action else {
1164 return None;
1165 };
1166
1167 Some(DueSessionPrompt {
1168 id: record.definition.id.clone(),
1169 name: record.definition.name.clone(),
1170 prompt: prompt.clone(),
1171 })
1172}
1173
1174fn next_scheduled_run(
1175 definition: &ScheduledTaskDefinition,
1176 last_base_run_at: Option<DateTime<Utc>>,
1177) -> Result<Option<NextScheduledRun>> {
1178 let Some(last_base_run_at) = last_base_run_at else {
1179 return Ok(None);
1180 };
1181
1182 let Some(next_base_run_at) = definition.schedule.next_base_fire_after(last_base_run_at)? else {
1183 return Ok(None);
1184 };
1185
1186 if definition.expires_at.is_some_and(|expiry| next_base_run_at > expiry) {
1187 return Ok(None);
1188 }
1189
1190 Ok(Some(NextScheduledRun {
1191 fire_at: definition.schedule.jittered_fire_at(&definition.id, next_base_run_at)?,
1192 base_fire_at: next_base_run_at,
1193 }))
1194}
1195
1196fn advance_record_runtime(
1197 record: &mut ScheduledTaskRecord,
1198 ran_at: DateTime<Utc>,
1199 status: TaskRunStatus,
1200) -> Result<bool> {
1201 record.runtime.last_run_at = Some(ran_at);
1202 record.runtime.last_status = Some(status);
1203
1204 match next_scheduled_run(&record.definition, record.runtime.next_base_run_at)? {
1205 Some(next_run) => {
1206 record.runtime.next_base_run_at = Some(next_run.base_fire_at);
1207 record.runtime.next_run_at = Some(next_run.fire_at);
1208 Ok(false)
1209 }
1210 None => {
1211 record.runtime.next_base_run_at = None;
1212 record.runtime.next_run_at = None;
1213 Ok(true)
1214 }
1215 }
1216}
1217
1218fn apply_run_outcome(record: &mut ScheduledTaskRecord, outcome: RunOutcome) -> Result<()> {
1219 let RunOutcome {
1220 ran_at,
1221 status,
1222 artifact_dir,
1223 events_file,
1224 last_message_file,
1225 } = outcome;
1226
1227 record.runtime.last_artifact_dir = artifact_dir;
1228 record.runtime.last_events_file = events_file;
1229 record.runtime.last_message_file = last_message_file;
1230
1231 advance_record_runtime(record, ran_at, status)?;
1232
1233 Ok(())
1234}
1235
1236#[derive(Debug, Clone)]
1237struct RunOutcome {
1238 ran_at: DateTime<Utc>,
1239 status: TaskRunStatus,
1240 artifact_dir: Option<PathBuf>,
1241 events_file: Option<PathBuf>,
1242 last_message_file: Option<PathBuf>,
1243}
1244
1245fn scheduled_workspace_label(definition: &ScheduledTaskDefinition) -> String {
1246 definition
1247 .workspace
1248 .as_ref()
1249 .map(|path| path.display().to_string())
1250 .unwrap_or_else(|| "<none>".to_string())
1251}
1252
1253fn resolve_scheduled_task_workspace(definition: &ScheduledTaskDefinition) -> Result<Option<PathBuf>> {
1254 definition
1255 .workspace
1256 .as_deref()
1257 .map(resolve_scheduled_workspace_path)
1258 .transpose()
1259 .with_context(|| {
1260 format!(
1261 "Failed to resolve scheduled task workspace {} for task {}",
1262 scheduled_workspace_label(definition),
1263 definition.id
1264 )
1265 })
1266}
1267
1268fn read_definition(path: &Path) -> Result<ScheduledTaskDefinition> {
1269 let raw = fs::read_to_string(path).with_context(|| format!("Failed to read {}", path.display()))?;
1270 toml::from_str(&raw).with_context(|| format!("Failed to parse {}", path.display()))
1271}
1272
1273fn atomic_write(path: &Path, content: &[u8]) -> Result<()> {
1274 VtCodePaths::write_private_file_atomic(path, content)
1275 .with_context(|| format!("Failed to atomically write {}", path.display()))
1276}
1277
1278fn try_acquire_claim(paths: &SchedulerPaths, id: &str) -> Result<bool> {
1279 paths.ensure_dirs()?;
1280 let path = paths.claim_path(id);
1281 match VtCodePaths::create_private_file(&path) {
1282 Ok(mut file) => {
1283 let timestamp = Utc::now().to_rfc3339();
1284 file.write_all(timestamp.as_bytes())
1285 .with_context(|| format!("Failed to write {}", path.display()))?;
1286 Ok(true)
1287 }
1288 Err(error)
1289 if error.chain().any(|cause| {
1290 cause
1291 .downcast_ref::<std::io::Error>()
1292 .is_some_and(|io_error| io_error.kind() == std::io::ErrorKind::AlreadyExists)
1293 }) =>
1294 {
1295 if claim_is_stale(&path)? {
1296 let _ = fs::remove_file(&path);
1297 return try_acquire_claim(paths, id);
1298 }
1299 Ok(false)
1300 }
1301 Err(error) => Err(error).with_context(|| format!("Failed to create {}", path.display())),
1302 }
1303}
1304
1305async fn try_acquire_claim_async(paths: SchedulerPaths, id: String) -> Result<bool> {
1306 tokio::task::spawn_blocking(move || try_acquire_claim(&paths, &id))
1307 .await
1308 .context("Scheduled task claim acquisition task panicked")?
1309}
1310
1311fn claim_is_stale(path: &Path) -> Result<bool> {
1312 let metadata = fs::metadata(path).with_context(|| format!("Failed to stat {}", path.display()))?;
1313 let modified = metadata
1314 .modified()
1315 .with_context(|| format!("Failed to read modification time for {}", path.display()))?;
1316 let elapsed = modified.elapsed().unwrap_or_default();
1317 Ok(elapsed >= Duration::from_secs(CLAIM_STALE_SECS))
1318}
1319
1320fn release_claim(paths: &SchedulerPaths, id: &str) -> Result<()> {
1321 let path = paths.claim_path(id);
1322 if path.exists() {
1323 fs::remove_file(&path).with_context(|| format!("Failed to remove {}", path.display()))?;
1324 }
1325 Ok(())
1326}
1327
1328async fn release_claim_async(paths: SchedulerPaths, id: String) -> Result<()> {
1329 tokio::task::spawn_blocking(move || release_claim(&paths, &id))
1330 .await
1331 .context("Scheduled task claim release task panicked")?
1332}
1333
1334fn summarize_task_name(summary: &str) -> String {
1335 let trimmed = summary.trim();
1336 if trimmed.is_empty() {
1337 return "Scheduled task".to_string();
1338 }
1339 let compact = vtcode_commons::formatting::collapse_whitespace(trimmed);
1340 let mut output = String::new();
1341 for ch in compact.chars().take(32) {
1342 output.push(ch);
1343 }
1344 output
1345}
1346
1347fn generate_task_id(name: &str, summary: &str, created_at: DateTime<Utc>) -> String {
1348 let counter = NEXT_TASK_COUNTER.fetch_add(1, Ordering::Relaxed);
1349 let seed = format!(
1350 "{name}|{summary}|{}|{}|{}",
1351 created_at.timestamp_nanos_opt().unwrap_or_default(),
1352 std::process::id(),
1353 counter
1354 );
1355 format!("{:08x}", stable_hash_u64(seed.as_bytes()) as u32)
1356}
1357
1358fn stable_hash_u64(bytes: &[u8]) -> u64 {
1359 let mut hash = 0xcbf29ce484222325u64;
1360 for byte in bytes {
1361 hash ^= u64::from(*byte);
1362 hash = hash.wrapping_mul(0x100000001b3);
1363 }
1364 hash
1365}
1366
1367fn humanize_interval(seconds: u64) -> String {
1368 match seconds {
1369 value if value % 86_400 == 0 => {
1370 let days = value / 86_400;
1371 if days == 1 {
1372 "every 1 day".to_string()
1373 } else {
1374 format!("every {days} days")
1375 }
1376 }
1377 value if value % 3_600 == 0 => {
1378 let hours = value / 3_600;
1379 if hours == 1 {
1380 "every 1 hour".to_string()
1381 } else {
1382 format!("every {hours} hours")
1383 }
1384 }
1385 value => {
1386 let minutes = value / 60;
1387 if minutes == 1 {
1388 "every 1 minute".to_string()
1389 } else {
1390 format!("every {minutes} minutes")
1391 }
1392 }
1393 }
1394}
1395
1396pub fn parse_local_datetime(raw: &str, now: DateTime<Local>) -> Result<DateTime<Utc>> {
1397 let trimmed = raw.trim();
1398 if trimmed.is_empty() {
1399 bail!("Time value cannot be empty");
1400 }
1401
1402 if let Ok(parsed) = DateTime::parse_from_rfc3339(trimmed) {
1403 return Ok(parsed.with_timezone(&Utc));
1404 }
1405
1406 for format in ["%Y-%m-%d %H:%M", "%Y-%m-%dT%H:%M", "%Y-%m-%d %H:%M:%S"] {
1407 if let Ok(naive) = NaiveDateTime::parse_from_str(trimmed, format) {
1408 return localize_naive_datetime(naive).map(|value| value.with_timezone(&Utc));
1409 }
1410 }
1411
1412 if let Some(captures) = TIME_ONLY_RE.captures(trimmed) {
1413 let mut hour = captures["hour"].parse::<u32>()?;
1414 let minute = captures
1415 .name("minute")
1416 .map(|value| value.as_str().parse::<u32>())
1417 .transpose()?
1418 .unwrap_or(0);
1419 if minute >= 60 {
1420 bail!("Invalid minute in time value");
1421 }
1422 if let Some(ampm) = captures.name("ampm").map(|value| value.as_str().to_ascii_lowercase()) {
1423 if hour == 0 || hour > 12 {
1424 bail!("Invalid 12-hour clock value");
1425 }
1426 if ampm == "pm" && hour != 12 {
1427 hour += 12;
1428 }
1429 if ampm == "am" && hour == 12 {
1430 hour = 0;
1431 }
1432 } else if hour >= 24 {
1433 bail!("Invalid 24-hour clock value");
1434 }
1435
1436 let time = NaiveTime::from_hms_opt(hour, minute, 0).ok_or_else(|| anyhow!("Invalid time value"))?;
1437 let today = now.date_naive();
1438 let naive = today.and_time(time);
1439 let mut localized = localize_naive_datetime(naive)?;
1440 if localized <= now {
1441 localized = localize_naive_datetime((today + ChronoDuration::days(1)).and_time(time))?;
1442 }
1443 return Ok(localized.with_timezone(&Utc));
1444 }
1445
1446 bail!("Unsupported time format. Use RFC3339, YYYY-MM-DD HH:MM, or a local time like 3pm")
1447}
1448
1449fn localize_naive_datetime(naive: NaiveDateTime) -> Result<DateTime<Local>> {
1450 match Local.from_local_datetime(&naive) {
1451 LocalResult::Single(value) => Ok(value),
1452 LocalResult::Ambiguous(first, _) => Ok(first),
1453 LocalResult::None => bail!("Local time does not exist due to timezone transition"),
1454 }
1455}
1456
1457fn xml_escape(value: String) -> String {
1458 value
1459 .replace('&', "&")
1460 .replace('<', "<")
1461 .replace('>', ">")
1462 .replace('"', """)
1463 .replace('\'', "'")
1464}
1465
1466#[derive(Debug, Clone)]
1467struct ParsedCron {
1468 minute: CronField,
1469 hour: CronField,
1470 day_of_month: CronField,
1471 month: CronField,
1472 day_of_week: CronField,
1473}
1474
1475impl ParsedCron {
1476 fn parse(expression: &str) -> Result<Self> {
1477 let parts = expression.split_whitespace().collect::<Vec<_>>();
1478 if parts.len() != 5 {
1479 bail!("Cron expressions require exactly 5 fields");
1480 }
1481
1482 Ok(Self {
1483 minute: CronField::parse(parts[0], 0, 59, false)?,
1484 hour: CronField::parse(parts[1], 0, 23, false)?,
1485 day_of_month: CronField::parse(parts[2], 1, 31, false)?,
1486 month: CronField::parse(parts[3], 1, 12, false)?,
1487 day_of_week: CronField::parse(parts[4], 0, 7, true)?,
1488 })
1489 }
1490
1491 fn next_after(&self, after: DateTime<Local>) -> Result<Option<DateTime<Local>>> {
1492 let mut candidate = after
1493 .with_second(0)
1494 .and_then(|value| value.with_nanosecond(0))
1495 .ok_or_else(|| anyhow!("Failed to normalize cron timestamp"))?
1496 + ChronoDuration::minutes(1);
1497 let horizon = candidate + ChronoDuration::days(366 * 5);
1498
1499 while candidate <= horizon {
1500 if self.matches(candidate) {
1501 return Ok(Some(candidate));
1502 }
1503 candidate += ChronoDuration::minutes(1);
1504 }
1505
1506 Ok(None)
1507 }
1508
1509 fn matches(&self, value: DateTime<Local>) -> bool {
1510 let month = value.month();
1511 let dom = value.day();
1512 let minute = value.minute();
1513 let hour = value.hour();
1514 let dow = value.weekday().num_days_from_sunday();
1515
1516 if !self.minute.contains(minute) || !self.hour.contains(hour) || !self.month.contains(month) {
1517 return false;
1518 }
1519
1520 let dom_matches = self.day_of_month.contains(dom);
1521 let dow_matches = self.day_of_week.contains(dow);
1522
1523 if self.day_of_month.is_wildcard && self.day_of_week.is_wildcard {
1524 return true;
1525 }
1526 if self.day_of_month.is_wildcard {
1527 return dow_matches;
1528 }
1529 if self.day_of_week.is_wildcard {
1530 return dom_matches;
1531 }
1532
1533 dom_matches || dow_matches
1534 }
1535}
1536
1537#[derive(Debug, Clone)]
1538struct CronField {
1539 values: BTreeSet<u32>,
1540 is_wildcard: bool,
1541}
1542
1543impl CronField {
1544 fn parse(raw: &str, min: u32, max: u32, is_day_of_week: bool) -> Result<Self> {
1545 if raw.contains(['L', 'W', '?']) {
1546 bail!("Unsupported cron syntax in field '{raw}'");
1547 }
1548 if raw.chars().any(|ch| ch.is_ascii_alphabetic()) {
1549 bail!("Named cron aliases are not supported in '{raw}'");
1550 }
1551
1552 let mut values = BTreeSet::new();
1553 let mut is_wildcard = false;
1554 for segment in raw.split(',') {
1555 let segment = segment.trim();
1556 if segment.is_empty() {
1557 bail!("Cron field contains an empty segment");
1558 }
1559 if segment == "*" {
1560 is_wildcard = true;
1561 values.extend(min..=max);
1562 continue;
1563 }
1564
1565 let (base, step) = if let Some((left, right)) = segment.split_once('/') {
1566 let step = right
1567 .parse::<u32>()
1568 .with_context(|| format!("Invalid step value in '{segment}'"))?;
1569 if step == 0 {
1570 bail!("Step value must be greater than zero");
1571 }
1572 (left, Some(step))
1573 } else {
1574 (segment, None)
1575 };
1576
1577 let mut base_values = if base == "*" {
1578 is_wildcard = true;
1579 (min..=max).collect::<Vec<_>>()
1580 } else if let Some((left, right)) = base.split_once('-') {
1581 let start = parse_cron_number(left, min, max, is_day_of_week)?;
1582 let end = parse_cron_number(right, min, max, is_day_of_week)?;
1583 if start > end {
1584 bail!("Invalid descending range '{base}'");
1585 }
1586 (start..=end).collect::<Vec<_>>()
1587 } else {
1588 let start = parse_cron_number(base, min, max, is_day_of_week)?;
1589 if let Some(step) = step {
1590 let mut values = Vec::new();
1591 let mut value = start;
1592 while value <= max {
1593 values.push(value);
1594 match value.checked_add(step) {
1595 Some(next) => value = next,
1596 None => break,
1597 }
1598 }
1599 values
1600 } else {
1601 vec![start]
1602 }
1603 };
1604
1605 if let Some(step) = step
1606 && (base == "*" || base.contains('-'))
1607 {
1608 let mut stepped = Vec::new();
1609 for (index, value) in base_values.iter().enumerate() {
1610 if (index as u32).is_multiple_of(step) {
1611 stepped.push(*value);
1612 }
1613 }
1614 base_values = stepped;
1615 }
1616
1617 values.extend(base_values);
1618 }
1619
1620 Ok(Self { values, is_wildcard })
1621 }
1622
1623 fn contains(&self, value: u32) -> bool {
1624 self.values.contains(&value)
1625 }
1626}
1627
1628fn parse_cron_number(raw: &str, min: u32, max: u32, is_day_of_week: bool) -> Result<u32> {
1629 let mut value = raw.parse::<u32>().with_context(|| format!("Invalid cron value '{raw}'"))?;
1630 if is_day_of_week && value == 7 {
1631 value = 0;
1632 }
1633 if !(min..=max).contains(&value) {
1634 bail!("Cron value '{raw}' is out of range");
1635 }
1636 Ok(value)
1637}