Skip to main content

ironflow_cli/commands/
schedule.rs

1//! Schedule subcommands: list, create, pause, resume, delete, trigger.
2
3use anyhow::Result;
4use clap::{Args, Subcommand, ValueEnum, value_parser};
5use comfy_table::{ContentArrangement, Table};
6use ironflow_sdk::IronflowClient;
7use ironflow_sdk::types::{CatchupPolicy, CreateScheduleRequest, OverlapPolicy, ScheduleResponse};
8use uuid::Uuid;
9
10use crate::confirm::confirm;
11use crate::output;
12
13/// Arguments for the `schedule` command group.
14#[derive(Debug, Args)]
15pub struct ScheduleArgs {
16    /// Schedule subcommand.
17    #[command(subcommand)]
18    pub command: ScheduleCommands,
19}
20
21/// Available schedule subcommands.
22#[derive(Debug, Subcommand)]
23pub enum ScheduleCommands {
24    /// List all schedules.
25    List,
26    /// Create a new schedule.
27    Create {
28        /// Workflow name.
29        workflow: String,
30        /// Cron expression (5 or 6 field format).
31        cron: String,
32        /// JSON inputs for the workflow (defaults to `{}`).
33        #[arg(long, default_value = "{}")]
34        inputs: String,
35        /// Queue priority, from -100 to 100, given to every run the schedule
36        /// creates. Defaults to the workflow priority.
37        #[arg(
38            long,
39            allow_negative_numbers = true,
40            value_parser = value_parser!(i16).range(-100..=100)
41        )]
42        priority: Option<i16>,
43        /// What the schedule does with the occurrences it missed while no
44        /// server fired it. Defaults to `latest`.
45        #[arg(long, value_enum)]
46        catchup: Option<CatchupArg>,
47        /// Most runs created to catch up under `--catchup all`, from 1 to
48        /// 1000. Defaults to 10.
49        #[arg(long, value_parser = value_parser!(i32).range(1..=1000))]
50        catchup_max: Option<i32>,
51        /// How far back, in seconds, a missed occurrence is still caught up,
52        /// from 60 to 2592000 (30 days). Defaults to 86400 (one day).
53        #[arg(
54            long = "catchup-window",
55            value_name = "SECONDS",
56            value_parser = value_parser!(i32).range(60..=2_592_000)
57        )]
58        catchup_window_secs: Option<i32>,
59        /// What the schedule does when an occurrence comes while one of its
60        /// runs is still active. Defaults to `allow`.
61        #[arg(long, value_enum)]
62        overlap: Option<OverlapArg>,
63        /// IANA timezone the cron expression is evaluated in, e.g.
64        /// `Europe/Paris`. Defaults to `UTC`.
65        #[arg(long, value_name = "IANA")]
66        timezone: Option<String>,
67    },
68    /// Pause a schedule (disable automatic triggers).
69    Pause {
70        /// Schedule ID.
71        id: Uuid,
72    },
73    /// Resume a paused schedule.
74    Resume {
75        /// Schedule ID.
76        id: Uuid,
77    },
78    /// Delete a schedule.
79    Delete {
80        /// Schedule ID.
81        id: Uuid,
82        /// Skip the interactive confirmation.
83        #[arg(long)]
84        yes: bool,
85    },
86    /// Trigger a schedule manually, creating a run immediately.
87    Trigger {
88        /// Schedule ID.
89        id: Uuid,
90    },
91}
92
93/// `--catchup` values.
94#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
95pub enum CatchupArg {
96    /// Run the most recent missed occurrence only.
97    Latest,
98    /// Run every missed occurrence, up to `--catchup-max`.
99    All,
100    /// Run no missed occurrence.
101    Skip,
102}
103
104impl From<CatchupArg> for CatchupPolicy {
105    fn from(arg: CatchupArg) -> Self {
106        match arg {
107            CatchupArg::Latest => Self::Latest,
108            CatchupArg::All => Self::All,
109            CatchupArg::Skip => Self::Skip,
110        }
111    }
112}
113
114/// `--overlap` values.
115#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
116pub enum OverlapArg {
117    /// Start another run even if one is still active.
118    Allow,
119    /// Drop the occurrence while a run of the schedule is still active.
120    Skip,
121}
122
123impl From<OverlapArg> for OverlapPolicy {
124    fn from(arg: OverlapArg) -> Self {
125        match arg {
126            OverlapArg::Allow => Self::Allow,
127            OverlapArg::Skip => Self::Skip,
128        }
129    }
130}
131
132/// The catch-up policy, with its bound under `all`.
133fn catchup_label(s: &ScheduleResponse) -> String {
134    match s.catchup {
135        CatchupPolicy::All => format!("all (max {})", s.catchup_max),
136        CatchupPolicy::Latest => "latest".to_string(),
137        CatchupPolicy::Skip => "skip".to_string(),
138    }
139}
140
141/// `active`, `paused` by a user, or `disabled: <reason>` when Ironflow
142/// disabled the schedule on an error.
143fn schedule_state(s: &ScheduleResponse) -> String {
144    match (&s.disabled_at, &s.last_error) {
145        (None, _) => "active".to_string(),
146        (Some(_), Some(error)) => format!("disabled: {error}"),
147        (Some(_), None) => "paused".to_string(),
148    }
149}
150
151fn schedules_table(schedules: &[ScheduleResponse]) -> Table {
152    let mut table = Table::new();
153    table.set_content_arrangement(ContentArrangement::Dynamic);
154    table.set_header(vec![
155        "ID",
156        "WORKFLOW",
157        "CRON",
158        "SOURCE",
159        "PRIORITY",
160        "TIMEZONE",
161        "CATCHUP",
162        "OVERLAP",
163        "ENABLED",
164        "NEXT TRIGGER",
165    ]);
166    for s in schedules {
167        let source = format!("{:?}", s.source).to_lowercase();
168        table.add_row(vec![
169            s.id.to_string(),
170            s.workflow_name.clone(),
171            s.cron_expression.clone(),
172            source.to_string(),
173            s.priority.to_string(),
174            s.timezone.clone(),
175            catchup_label(s),
176            format!("{:?}", s.overlap).to_lowercase(),
177            schedule_state(s),
178            s.next_trigger_at
179                .as_ref()
180                .map(|d| d.to_string())
181                .unwrap_or_else(|| "-".to_string()),
182        ]);
183    }
184    table
185}
186
187/// Execute a schedule subcommand.
188///
189/// # Errors
190///
191/// Returns an error on API failure, invalid JSON inputs, or an unconfirmed
192/// destructive command.
193pub async fn execute(client: &IronflowClient, args: &ScheduleArgs, json_mode: bool) -> Result<()> {
194    match &args.command {
195        ScheduleCommands::List => {
196            let response = client.list_schedules().await?;
197            if json_mode {
198                output::print_json(&response)?;
199            } else {
200                println!("{}", schedules_table(&response.data));
201            }
202            Ok(())
203        }
204        ScheduleCommands::Create {
205            workflow,
206            cron,
207            inputs,
208            priority,
209            catchup,
210            catchup_max,
211            catchup_window_secs,
212            overlap,
213            timezone,
214        } => {
215            let parsed_inputs: serde_json::Value =
216                serde_json::from_str(inputs).map_err(|e| anyhow::anyhow!("invalid JSON: {e}"))?;
217            let response = client
218                .create_schedule(&CreateScheduleRequest {
219                    workflow_name: workflow.clone(),
220                    cron_expression: cron.clone(),
221                    inputs: Some(parsed_inputs),
222                    priority: priority.map(i32::from),
223                    catchup: catchup.map(CatchupPolicy::from),
224                    catchup_max: *catchup_max,
225                    catchup_window_secs: *catchup_window_secs,
226                    overlap: overlap.map(OverlapPolicy::from),
227                    timezone: timezone.clone(),
228                })
229                .await?;
230            if json_mode {
231                output::print_json(&response)?;
232            } else {
233                println!("Schedule {} created", response.data.id);
234            }
235            Ok(())
236        }
237        ScheduleCommands::Pause { id } => {
238            let response = client.pause_schedule(*id).await?;
239            if json_mode {
240                output::print_json(&response)?;
241            } else {
242                println!("Schedule {} paused", response.data.id);
243            }
244            Ok(())
245        }
246        ScheduleCommands::Resume { id } => {
247            let response = client.resume_schedule(*id).await?;
248            if json_mode {
249                output::print_json(&response)?;
250            } else {
251                println!("Schedule {} resumed", response.data.id);
252            }
253            Ok(())
254        }
255        ScheduleCommands::Delete { id, yes } => {
256            let prompt = format!("Delete schedule {id}?");
257            confirm(&prompt, *yes)?;
258            client.delete_schedule(*id).await?;
259            if json_mode {
260                output::print_json(&serde_json::json!({"deleted": id.to_string()}))?;
261            } else {
262                println!("Schedule {id} deleted");
263            }
264            Ok(())
265        }
266        ScheduleCommands::Trigger { id } => {
267            let response = client.trigger_schedule(*id).await?;
268            if json_mode {
269                output::print_json(&response)?;
270            } else {
271                println!("Schedule {} triggered", response.data.id);
272            }
273            Ok(())
274        }
275    }
276}
277
278#[cfg(test)]
279mod tests {
280    use serde_json::{from_value, json};
281
282    use super::*;
283
284    fn schedule(disabled_at: Option<&str>, last_error: Option<&str>) -> ScheduleResponse {
285        from_value(json!({
286            "id": "01a10fdb-f467-7982-826b-0c99470c7264",
287            "workflow_name": "deploy",
288            "cron_expression": "0 0 30 2 *",
289            "inputs": {},
290            "source": "api",
291            "priority": -10,
292            "catchup": "latest",
293            "catchup_max": 10,
294            "catchup_window_secs": 86400,
295            "overlap": "allow",
296            "timezone": "UTC",
297            "disabled_at": disabled_at,
298            "last_triggered_at": null,
299            "next_trigger_at": null,
300            "last_error": last_error,
301            "created_by_user_id": null,
302            "created_at": "2026-10-06T08:00:00Z",
303            "updated_at": "2026-10-06T08:00:00Z"
304        }))
305        .unwrap()
306    }
307
308    #[test]
309    fn schedule_state_shows_why_ironflow_disabled_a_schedule() {
310        let s = schedule(
311            Some("2026-10-06T08:00:00Z"),
312            Some("cannot compute next trigger"),
313        );
314        assert_eq!(schedule_state(&s), "disabled: cannot compute next trigger");
315        assert!(
316            schedules_table(&[s])
317                .to_string()
318                .contains("disabled: cannot compute next trigger")
319        );
320    }
321
322    #[test]
323    fn schedule_state_tells_a_user_pause_from_an_active_schedule() {
324        assert_eq!(
325            schedule_state(&schedule(Some("2026-10-06T08:00:00Z"), None)),
326            "paused"
327        );
328        assert_eq!(schedule_state(&schedule(None, None)), "active");
329    }
330
331    #[test]
332    fn schedules_table_shows_the_priority() {
333        let output = schedules_table(&[schedule(None, None)]).to_string();
334        assert!(
335            output.contains("PRIORITY"),
336            "header missing from:\n{output}"
337        );
338        assert!(output.contains("-10"), "priority missing from:\n{output}");
339    }
340
341    #[test]
342    fn schedules_table_shows_timezone_and_policies() {
343        let mut s = schedule(None, None);
344        s.timezone = "Europe/Paris".to_string();
345        s.catchup = CatchupPolicy::All;
346        s.catchup_max = 24;
347        s.overlap = OverlapPolicy::Skip;
348
349        let output = schedules_table(&[s]).to_string();
350
351        for header in ["TIMEZONE", "CATCHUP", "OVERLAP"] {
352            assert!(output.contains(header), "{header} missing from:\n{output}");
353        }
354        assert!(
355            output.contains("Europe/Paris"),
356            "timezone missing from:\n{output}"
357        );
358        assert!(
359            output.contains("all (max 24)"),
360            "catchup missing from:\n{output}"
361        );
362        assert!(output.contains("skip"), "overlap missing from:\n{output}");
363    }
364
365    #[test]
366    fn catchup_label_names_the_policy() {
367        let mut s = schedule(None, None);
368        assert_eq!(catchup_label(&s), "latest");
369        s.catchup = CatchupPolicy::Skip;
370        assert_eq!(catchup_label(&s), "skip");
371    }
372}