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