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