Skip to main content

vtcode_core/scheduler/
mod.rs

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('&', "&amp;")
1460        .replace('<', "&lt;")
1461        .replace('>', "&gt;")
1462        .replace('"', "&quot;")
1463        .replace('\'', "&apos;")
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}