Skip to main content

ironflow_cli/commands/
run.rs

1//! Run subcommands: create, list, get, cancel, approve, reject, input, reject-input, retry,
2//! replay, watch, plan, diff.
3
4use std::fs;
5use std::io::{Write as _, stdout};
6use std::path::PathBuf;
7use std::slice;
8use std::time::Duration;
9
10use anyhow::{Context, Result, anyhow, bail};
11use clap::{Args, Subcommand, value_parser};
12use futures_util::StreamExt;
13use humantime::format_duration;
14use ironflow_sdk::IronflowClient;
15use ironflow_sdk::client::ListRunsFilter;
16use ironflow_sdk::types::{ConcurrencyLimit, CreateRunRequest, PlanWorkflowRequest, RunStatus};
17use ironflow_types::parse_concurrency_limit as shared_parse_concurrency_limit;
18use serde_json::{Map, Value, from_str, json, to_string};
19use tokio::time::timeout as tokio_timeout;
20use uuid::Uuid;
21
22use crate::output;
23
24/// Arguments for the `run` command group.
25#[derive(Debug, Args)]
26pub struct RunArgs {
27    /// Run subcommand.
28    #[command(subcommand)]
29    pub command: RunCommands,
30}
31
32/// Available run subcommands.
33#[derive(Debug, Subcommand)]
34pub enum RunCommands {
35    /// Create a new run for a workflow.
36    Create {
37        /// Workflow name to trigger.
38        workflow: String,
39        /// JSON payload (inline string).
40        #[arg(long, group = "payload_source")]
41        payload: Option<String>,
42        /// Path to a JSON file containing the payload.
43        #[arg(long, group = "payload_source")]
44        payload_file: Option<PathBuf>,
45        /// How many times to replay the run automatically after a transient
46        /// failure. Defaults to 0 (no automatic retry).
47        #[arg(long)]
48        max_retries: Option<u32>,
49        /// Idempotency key making the call safe to replay.
50        ///
51        /// Reusing the same key returns the run it already created instead of
52        /// starting a second one. Valid for 24 hours. At most 255 printable
53        /// ASCII characters.
54        #[arg(long)]
55        idempotency_key: Option<String>,
56        /// Maximum cumulative cost for this run, in USD. Overrides the
57        /// workflow and server defaults.
58        #[arg(long = "max-cost", value_name = "USD")]
59        max_cost: Option<f64>,
60        /// Exclusivity key: the call is refused while a non-terminal run holds
61        /// the same key. Released once that run completes, fails or is
62        /// cancelled. At most 255 bytes.
63        #[arg(long)]
64        concurrency_key: Option<String>,
65        /// Concurrency group the run joins, as `GROUP=N`: a worker only starts
66        /// the run while fewer than N root runs of GROUP are running. Repeat
67        /// the flag to join several groups.
68        #[arg(
69            long = "concurrency-limit",
70            value_name = "GROUP=N",
71            value_parser = parse_concurrency_limit
72        )]
73        concurrency_limits: Vec<ConcurrencyLimit>,
74        /// Queue priority, from -100 to 100. Workers pick the pending run with
75        /// the highest priority first. Defaults to the workflow priority, `0`
76        /// unless the handler declares one.
77        #[arg(
78            long,
79            allow_negative_numbers = true,
80            value_parser = value_parser!(i16).range(-100..=100)
81        )]
82        priority: Option<i16>,
83        /// Tag a worker must carry to take the run, added to the ones the
84        /// workflow requires. Repeat the flag to require several tags.
85        #[arg(long = "worker-tag", value_name = "TAG")]
86        worker_tags: Vec<String>,
87    },
88    /// List runs with optional filters.
89    List {
90        /// Filter by run status (pending, running, completed, failed, etc.).
91        #[arg(long)]
92        status: Option<String>,
93        /// Filter by workflow name.
94        #[arg(long)]
95        workflow: Option<String>,
96        /// Filter by author: the user ID that triggered the run.
97        ///
98        /// Also matches runs triggered by one of that user's API keys.
99        #[arg(long)]
100        created_by: Option<Uuid>,
101        /// Filter by concurrency group: only runs that belong to this group.
102        #[arg(long)]
103        concurrency_group: Option<String>,
104        /// Filter by priority: only runs with exactly this priority.
105        #[arg(
106            long,
107            allow_negative_numbers = true,
108            value_parser = value_parser!(i16).range(-100..=100)
109        )]
110        priority: Option<i16>,
111        /// Page number (1-based).
112        #[arg(long)]
113        page: Option<u32>,
114        /// Items per page.
115        #[arg(long)]
116        per_page: Option<u32>,
117    },
118    /// Get details of a specific run.
119    Get {
120        /// Run UUID.
121        id: Uuid,
122    },
123    /// Cancel a pending or running run.
124    Cancel {
125        /// Run UUID.
126        id: Uuid,
127    },
128    /// Approve a run waiting for approval.
129    Approve {
130        /// Run UUID.
131        id: Uuid,
132    },
133    /// Reject a run waiting for approval, failing it.
134    Reject {
135        /// Run UUID.
136        id: Uuid,
137    },
138    /// Answer a human input step; the value must match the step's JSON schema.
139    Input {
140        /// Run UUID.
141        id: Uuid,
142        /// Human input step UUID.
143        step_id: Uuid,
144        /// JSON answer (inline string).
145        #[arg(long, group = "value_source")]
146        value: Option<String>,
147        /// Path to a JSON file containing the answer.
148        #[arg(long, group = "value_source")]
149        value_file: Option<PathBuf>,
150    },
151    /// Reject a human input step; the handler receives the rejection.
152    RejectInput {
153        /// Run UUID.
154        id: Uuid,
155        /// Human input step UUID.
156        step_id: Uuid,
157        /// Why the input is refused.
158        #[arg(long)]
159        reason: Option<String>,
160    },
161    /// Retry a failed run.
162    Retry {
163        /// Run UUID.
164        id: Uuid,
165        /// Force retry even when the handler version has changed since the
166        /// original run.
167        #[arg(long)]
168        force: bool,
169    },
170    /// Replay a finished run on the current handler version.
171    Replay {
172        /// Run UUID.
173        id: Uuid,
174    },
175    /// Watch a run in real time via SSE.
176    Watch {
177        /// Run UUID.
178        id: Uuid,
179        /// Only show step transitions, not log data.
180        #[arg(long)]
181        no_logs: bool,
182        /// Stop watching after this duration (e.g. "30s", "5m", "1h").
183        #[arg(long, value_parser = parse_humantime)]
184        timeout: Option<Duration>,
185    },
186    /// Show the execution plan for a workflow without running it.
187    Plan {
188        /// Workflow name.
189        workflow: String,
190        /// JSON input (inline string).
191        #[arg(long, group = "plan_input_source")]
192        input: Option<String>,
193        /// Path to a JSON file containing the input.
194        #[arg(long, group = "plan_input_source")]
195        input_file: Option<PathBuf>,
196        /// How deep sub-workflows are expanded (default 3).
197        #[arg(long)]
198        max_depth: Option<u32>,
199        /// Skip duration estimation from run history.
200        #[arg(long)]
201        no_estimates: bool,
202    },
203    /// Compare two runs of the same workflow side by side.
204    Diff {
205        /// First run UUID.
206        run_a: Uuid,
207        /// Second run UUID.
208        run_b: Uuid,
209    },
210}
211
212/// Parse a human-readable duration string (e.g. "30s", "5m", "1h").
213fn parse_humantime(s: &str) -> Result<Duration, String> {
214    humantime::parse_duration(s).map_err(|e| e.to_string())
215}
216
217/// Parse a `GROUP=N` concurrency limit into the SDK type.
218///
219/// Delegates to [`ironflow_types::parse_concurrency_limit`]; the group and the
220/// limit are validated by the API.
221fn parse_concurrency_limit(s: &str) -> Result<ConcurrencyLimit, String> {
222    let (group, limit) = shared_parse_concurrency_limit(s)?;
223    let limit =
224        i32::try_from(limit).map_err(|e| format!("invalid limit '{limit}' in '{s}': {e}"))?;
225    Ok(ConcurrencyLimit { group, limit })
226}
227
228/// Terminal event types that signal the run is done.
229const TERMINAL_EVENTS: &[&str] = &["run_completed", "run_failed", "run_cancelled"];
230
231/// Resolve the payload from inline string or file.
232fn resolve_payload(payload: Option<&str>, payload_file: Option<&PathBuf>) -> Result<Value> {
233    match (payload, payload_file) {
234        (Some(raw), _) => from_str(raw).context("invalid JSON in --payload"),
235        (_, Some(path)) => {
236            let content = fs::read_to_string(path)
237                .with_context(|| format!("cannot read payload file: {}", path.display()))?;
238            from_str(&content).with_context(|| format!("invalid JSON in {}", path.display()))
239        }
240        (None, None) => Ok(Value::Object(Map::new())),
241    }
242}
243
244/// Reject a `--max-cost` value the API would refuse anyway.
245///
246/// Catching it client-side turns a 400 round-trip into an immediate, readable
247/// error.
248///
249/// # Errors
250///
251/// Returns an error when the value is negative or not a finite number.
252fn validate_max_cost(max_cost: Option<f64>) -> Result<()> {
253    match max_cost {
254        Some(value) if !value.is_finite() => {
255            anyhow::bail!("--max-cost must be a finite number, got {value}")
256        }
257        Some(value) if value < 0.0 => {
258            anyhow::bail!("--max-cost must be zero or positive, got {value}")
259        }
260        _ => Ok(()),
261    }
262}
263
264/// Execute a run subcommand.
265///
266/// # Errors
267///
268/// Returns an error on API failure or invalid input.
269pub async fn execute(
270    client: &IronflowClient,
271    args: &RunArgs,
272    json_mode: bool,
273    _verbose: bool,
274) -> Result<()> {
275    match &args.command {
276        RunCommands::Create {
277            workflow,
278            payload,
279            payload_file,
280            max_retries,
281            idempotency_key,
282            max_cost,
283            concurrency_key,
284            concurrency_limits,
285            priority,
286            worker_tags,
287        } => {
288            validate_max_cost(*max_cost)?;
289            let payload_value = resolve_payload(payload.as_deref(), payload_file.as_ref())?;
290            let payload_map = payload_value
291                .as_object()
292                .context("payload must be a JSON object")?
293                .clone();
294            let request: CreateRunRequest = CreateRunRequest::builder()
295                .workflow(workflow.clone())
296                .payload(Some(payload_map))
297                // The generated SDK models the field as i32; the API rejects
298                // anything negative, and clap already refuses it here.
299                .max_retries(max_retries.map(|n| n as i32))
300                .max_cost_usd(*max_cost)
301                .concurrency_key(concurrency_key.clone())
302                .concurrency_limits(concurrency_limits.clone())
303                .priority(priority.map(i32::from))
304                .worker_tags(worker_tags.clone())
305                .try_into()
306                .context("failed to build CreateRunRequest")?;
307
308            let response = match idempotency_key {
309                Some(key) => client.create_run_idempotent(&request, key).await?,
310                None => client.create_run(&request).await?,
311            };
312            output::print_output(json_mode, &response, || {
313                output::runs_table(slice::from_ref(&response.data))
314            })?;
315        }
316        RunCommands::List {
317            status,
318            workflow,
319            created_by,
320            concurrency_group,
321            priority,
322            page,
323            per_page,
324        } => {
325            let filter = ListRunsFilter {
326                status: status.as_deref(),
327                workflow: workflow.as_deref(),
328                created_by: *created_by,
329                concurrency_group: concurrency_group.as_deref(),
330                priority: *priority,
331                page: *page,
332                per_page: *per_page,
333                ..Default::default()
334            };
335            let response = client.list_runs_filtered(&filter).await?;
336            output::print_output(json_mode, &response, || output::runs_table(&response.data))?;
337        }
338        RunCommands::Get { id } => {
339            let response = client.get_run(*id).await?;
340            output::print_output(json_mode, &response, || {
341                output::run_detail_table(&response.data)
342            })?;
343
344            if !json_mode && !response.data.steps.is_empty() {
345                let mut out = stdout().lock();
346                writeln!(out)?;
347                writeln!(out, "Steps:")?;
348                writeln!(out, "{}", output::steps_table(&response.data.steps))?;
349            }
350        }
351        RunCommands::Cancel { id } => {
352            let response = client.cancel_run(*id).await?;
353            output::print_output(json_mode, &response, || {
354                output::cancelled_table(&response.data)
355            })?;
356        }
357        RunCommands::Approve { id } => {
358            let response = client.approve_run(*id).await?;
359            output::print_output(json_mode, &response, || {
360                output::runs_table(slice::from_ref(&response.data))
361            })?;
362            // A multi-approver gate stays open until enough distinct users voted.
363            if !json_mode && matches!(response.data.status, RunStatus::AwaitingApproval) {
364                println!("Approval recorded; more approvals are required.");
365            }
366        }
367        RunCommands::Reject { id } => {
368            let response = client.reject_run(*id).await?;
369            output::print_output(json_mode, &response, || {
370                output::runs_table(slice::from_ref(&response.data))
371            })?;
372        }
373        RunCommands::Input {
374            id,
375            step_id,
376            value,
377            value_file,
378        } => {
379            let answer = resolve_payload(value.as_deref(), value_file.as_ref())?;
380            let response = client.submit_input(*id, *step_id, &answer).await?;
381            output::print_output(json_mode, &response, || {
382                output::runs_table(slice::from_ref(&response.data))
383            })?;
384        }
385        RunCommands::RejectInput {
386            id,
387            step_id,
388            reason,
389        } => {
390            let response = client
391                .reject_input(*id, *step_id, reason.as_deref())
392                .await?;
393            output::print_output(json_mode, &response, || {
394                output::runs_table(slice::from_ref(&response.data))
395            })?;
396        }
397        RunCommands::Retry { id, force } => {
398            let response = client.retry_run(*id, *force).await?;
399            output::print_output(json_mode, &response, || {
400                output::runs_table(slice::from_ref(&response.data))
401            })?;
402        }
403        RunCommands::Replay { id } => {
404            let response = client.replay_run(*id).await?;
405            output::print_output(json_mode, &response, || {
406                output::runs_table(slice::from_ref(&response.data))
407            })?;
408        }
409        RunCommands::Watch {
410            id,
411            no_logs,
412            timeout,
413        } => {
414            execute_watch(client, *id, *no_logs, *timeout, json_mode).await?;
415        }
416        RunCommands::Plan {
417            workflow,
418            input,
419            input_file,
420            max_depth,
421            no_estimates,
422        } => {
423            let payload = resolve_payload(input.as_deref(), input_file.as_ref())?;
424            let payload_map = payload
425                .as_object()
426                .context("input must be a JSON object")?
427                .clone();
428            let request: PlanWorkflowRequest = PlanWorkflowRequest::builder()
429                .payload(Some(payload_map))
430                // The generated SDK models the depth as i32; the API rejects
431                // anything outside 1..=10.
432                .max_depth(max_depth.map(|d| d as i32))
433                .estimate_durations(Some(!*no_estimates))
434                .try_into()
435                .context("failed to build PlanWorkflowRequest")?;
436            let response = client.plan_workflow(workflow, &request).await?;
437            output::render_execution_plan(&mut stdout().lock(), json_mode, &response)?;
438        }
439        RunCommands::Diff { run_a, run_b } => {
440            execute_diff(client, *run_a, *run_b, json_mode).await?;
441        }
442    }
443    Ok(())
444}
445
446/// Watch a run in real time via SSE.
447async fn execute_watch(
448    client: &IronflowClient,
449    run_id: Uuid,
450    no_logs: bool,
451    timeout: Option<Duration>,
452    json_mode: bool,
453) -> Result<()> {
454    let run = client.get_run(run_id).await?;
455    let status = run.data.run.status;
456    if matches!(
457        status,
458        RunStatus::Completed | RunStatus::Failed | RunStatus::Cancelled
459    ) {
460        if json_mode {
461            output::print_output(json_mode, &run, || output::run_detail_table(&run.data))?;
462        } else {
463            let mut out = stdout().lock();
464            writeln!(out, "Run {run_id} already in terminal state: {status}")?;
465        }
466        return Ok(());
467    }
468
469    let watch_fut = async {
470        let mut stream = client.events(Some(run_id), None).await?;
471        let mut out = stdout().lock();
472
473        while let Some(event) = stream.next().await {
474            match event {
475                Ok(ev) => {
476                    if no_logs
477                        && !ev.event_type.starts_with("run_")
478                        && !ev.event_type.starts_with("step_")
479                    {
480                        continue;
481                    }
482
483                    if json_mode {
484                        let obj = json!({
485                            "event": ev.event_type,
486                            "data": ev.data,
487                        });
488                        writeln!(out, "{}", to_string(&obj)?)?;
489                    } else {
490                        writeln!(out, "[{}] {}", ev.event_type, ev.data)?;
491                    }
492
493                    if TERMINAL_EVENTS.contains(&ev.event_type.as_str()) {
494                        break;
495                    }
496                }
497                Err(e) => {
498                    return Err(anyhow!("SSE stream error: {e}"));
499                }
500            }
501        }
502
503        Ok::<(), anyhow::Error>(())
504    };
505
506    match timeout {
507        Some(dur) => {
508            tokio_timeout(dur, watch_fut).await.unwrap_or_else(|_| {
509                eprintln!("Timeout reached after {}", format_duration(dur));
510                Ok(())
511            })?;
512        }
513        None => {
514            watch_fut.await?;
515        }
516    }
517
518    Ok(())
519}
520
521/// Compare two runs of the same workflow side by side.
522async fn execute_diff(
523    client: &IronflowClient,
524    run_a_id: Uuid,
525    run_b_id: Uuid,
526    json_mode: bool,
527) -> Result<()> {
528    if run_a_id == run_b_id {
529        bail!("both run IDs are the same; nothing to diff");
530    }
531
532    let (a, b) = tokio::try_join!(client.get_run(run_a_id), client.get_run(run_b_id))?;
533
534    if a.data.run.workflow_name != b.data.run.workflow_name {
535        bail!(
536            "cannot diff runs from different workflows: '{}' vs '{}'",
537            a.data.run.workflow_name,
538            b.data.run.workflow_name
539        );
540    }
541
542    if json_mode {
543        let diff = json!({
544            "run_a": a.data,
545            "run_b": b.data,
546        });
547        output::print_json(&diff)?;
548    } else {
549        let table = output::run_diff_table(&a.data, &b.data);
550        let mut out = stdout().lock();
551        writeln!(out, "{table}")?;
552    }
553
554    Ok(())
555}
556
557#[cfg(test)]
558mod tests {
559    use std::io::Write;
560
561    use tempfile::NamedTempFile;
562
563    use super::*;
564
565    #[test]
566    fn parse_concurrency_limit_reads_group_and_limit() {
567        let limit = parse_concurrency_limit("repo:acme=2").unwrap();
568        assert_eq!(limit.group, "repo:acme");
569        assert_eq!(limit.limit, 2);
570    }
571
572    #[test]
573    fn parse_concurrency_limit_rejects_a_limit_above_i32() {
574        assert!(parse_concurrency_limit("repo:acme=4294967295").is_err());
575    }
576
577    #[test]
578    fn parse_concurrency_limit_propagates_shape_errors() {
579        assert!(parse_concurrency_limit("repo:acme").is_err());
580    }
581
582    #[test]
583    fn resolve_payload_none_returns_empty_object() {
584        let value = resolve_payload(None, None).unwrap();
585        assert!(value.is_object());
586        assert!(value.as_object().unwrap().is_empty());
587    }
588
589    #[test]
590    fn resolve_payload_inline_valid_json() {
591        let value = resolve_payload(Some(r#"{"key": "value"}"#), None).unwrap();
592        assert_eq!(value["key"], "value");
593    }
594
595    #[test]
596    fn resolve_payload_inline_invalid_json() {
597        let result = resolve_payload(Some("not json"), None);
598        assert!(result.is_err());
599        assert!(result.unwrap_err().to_string().contains("invalid JSON"));
600    }
601
602    #[test]
603    fn resolve_payload_file_valid() {
604        let mut tmp = NamedTempFile::new().unwrap();
605        write!(tmp, r#"{{"workflow": "test"}}"#).unwrap();
606        let path = tmp.path().to_path_buf();
607
608        let value = resolve_payload(None, Some(&path)).unwrap();
609        assert_eq!(value["workflow"], "test");
610    }
611
612    #[test]
613    fn resolve_payload_file_not_found() {
614        let path = PathBuf::from("/nonexistent/payload.json");
615        let result = resolve_payload(None, Some(&path));
616        assert!(result.is_err());
617        assert!(result.unwrap_err().to_string().contains("cannot read"));
618    }
619
620    #[test]
621    fn resolve_payload_file_invalid_json() {
622        let mut tmp = NamedTempFile::new().unwrap();
623        write!(tmp, "not valid json").unwrap();
624        let path = tmp.path().to_path_buf();
625
626        let result = resolve_payload(None, Some(&path));
627        assert!(result.is_err());
628        assert!(result.unwrap_err().to_string().contains("invalid JSON"));
629    }
630}