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