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