1use 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#[derive(Debug, Args)]
15pub struct ScheduleArgs {
16 #[command(subcommand)]
18 pub command: ScheduleCommands,
19}
20
21#[derive(Debug, Subcommand)]
23pub enum ScheduleCommands {
24 List,
26 Create {
28 workflow: String,
30 cron: String,
32 #[arg(long, default_value = "{}")]
34 inputs: String,
35 #[arg(
38 long,
39 allow_negative_numbers = true,
40 value_parser = value_parser!(i16).range(-100..=100)
41 )]
42 priority: Option<i16>,
43 #[arg(long, value_enum)]
46 catchup: Option<CatchupArg>,
47 #[arg(long, value_parser = value_parser!(i32).range(1..=1000))]
50 catchup_max: Option<i32>,
51 #[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 #[arg(long, value_enum)]
62 overlap: Option<OverlapArg>,
63 #[arg(long, value_name = "IANA")]
66 timezone: Option<String>,
67 },
68 Pause {
70 id: Uuid,
72 },
73 Resume {
75 id: Uuid,
77 },
78 Delete {
80 id: Uuid,
82 #[arg(long)]
84 yes: bool,
85 },
86 Trigger {
88 id: Uuid,
90 },
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
95pub enum CatchupArg {
96 Latest,
98 All,
100 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#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
116pub enum OverlapArg {
117 Allow,
119 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
132fn 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
141fn 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
187pub 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}