Skip to main content

ironflow_cli/commands/
run.rs

1//! Run subcommands: create, list, get, cancel, approve, retry.
2
3use std::fs;
4use std::io::Write as _;
5use std::path::PathBuf;
6use std::slice;
7
8use anyhow::{Context, Result};
9use clap::{Args, Subcommand};
10use ironflow_sdk::IronflowClient;
11use ironflow_sdk::client::ListRunsFilter;
12use ironflow_sdk::types::CreateRunRequest;
13use uuid::Uuid;
14
15use crate::output;
16
17/// Arguments for the `run` command group.
18#[derive(Debug, Args)]
19pub struct RunArgs {
20    /// Run subcommand.
21    #[command(subcommand)]
22    pub command: RunCommands,
23}
24
25/// Available run subcommands.
26#[derive(Debug, Subcommand)]
27pub enum RunCommands {
28    /// Create a new run for a workflow.
29    Create {
30        /// Workflow name to trigger.
31        workflow: String,
32        /// JSON payload (inline string).
33        #[arg(long, group = "payload_source")]
34        payload: Option<String>,
35        /// Path to a JSON file containing the payload.
36        #[arg(long, group = "payload_source")]
37        payload_file: Option<PathBuf>,
38        /// Idempotency key making the call safe to replay.
39        ///
40        /// Reusing the same key returns the run it already created instead of
41        /// starting a second one. Valid for 24 hours. At most 255 printable
42        /// ASCII characters.
43        #[arg(long)]
44        idempotency_key: Option<String>,
45        /// Maximum cumulative cost for this run, in USD. Overrides the
46        /// workflow and server defaults.
47        #[arg(long = "max-cost", value_name = "USD")]
48        max_cost: Option<f64>,
49    },
50    /// List runs with optional filters.
51    List {
52        /// Filter by run status (pending, running, completed, failed, etc.).
53        #[arg(long)]
54        status: Option<String>,
55        /// Filter by workflow name.
56        #[arg(long)]
57        workflow: Option<String>,
58        /// Filter by author: the user ID that triggered the run.
59        ///
60        /// Also matches runs triggered by one of that user's API keys.
61        #[arg(long)]
62        created_by: Option<Uuid>,
63        /// Page number (1-based).
64        #[arg(long)]
65        page: Option<u32>,
66        /// Items per page.
67        #[arg(long)]
68        per_page: Option<u32>,
69    },
70    /// Get details of a specific run.
71    Get {
72        /// Run UUID.
73        id: Uuid,
74    },
75    /// Cancel a pending or running run.
76    Cancel {
77        /// Run UUID.
78        id: Uuid,
79    },
80    /// Approve a run waiting for approval.
81    Approve {
82        /// Run UUID.
83        id: Uuid,
84    },
85    /// Retry a failed run.
86    Retry {
87        /// Run UUID.
88        id: Uuid,
89    },
90}
91
92/// Resolve the payload from inline string or file.
93fn resolve_payload(
94    payload: Option<&str>,
95    payload_file: Option<&PathBuf>,
96) -> Result<serde_json::Value> {
97    match (payload, payload_file) {
98        (Some(raw), _) => serde_json::from_str(raw).context("invalid JSON in --payload"),
99        (_, Some(path)) => {
100            let content = fs::read_to_string(path)
101                .with_context(|| format!("cannot read payload file: {}", path.display()))?;
102            serde_json::from_str(&content)
103                .with_context(|| format!("invalid JSON in {}", path.display()))
104        }
105        (None, None) => Ok(serde_json::Value::Object(serde_json::Map::new())),
106    }
107}
108
109/// Reject a `--max-cost` value the API would refuse anyway.
110///
111/// Catching it client-side turns a 400 round-trip into an immediate, readable
112/// error.
113///
114/// # Errors
115///
116/// Returns an error when the value is negative or not a finite number.
117fn validate_max_cost(max_cost: Option<f64>) -> Result<()> {
118    match max_cost {
119        Some(value) if !value.is_finite() => {
120            anyhow::bail!("--max-cost must be a finite number, got {value}")
121        }
122        Some(value) if value < 0.0 => {
123            anyhow::bail!("--max-cost must be zero or positive, got {value}")
124        }
125        _ => Ok(()),
126    }
127}
128
129/// Execute a run subcommand.
130///
131/// # Errors
132///
133/// Returns an error on API failure or invalid input.
134pub async fn execute(
135    client: &IronflowClient,
136    args: &RunArgs,
137    json_mode: bool,
138    _verbose: bool,
139) -> Result<()> {
140    match &args.command {
141        RunCommands::Create {
142            workflow,
143            payload,
144            payload_file,
145            idempotency_key,
146            max_cost,
147        } => {
148            validate_max_cost(*max_cost)?;
149            let payload_value = resolve_payload(payload.as_deref(), payload_file.as_ref())?;
150            let payload_map = payload_value
151                .as_object()
152                .context("payload must be a JSON object")?
153                .clone();
154            let request: CreateRunRequest = CreateRunRequest::builder()
155                .workflow(workflow.clone())
156                .payload(Some(payload_map))
157                .max_cost_usd(*max_cost)
158                .try_into()
159                .context("failed to build CreateRunRequest")?;
160
161            let response = match idempotency_key {
162                Some(key) => client.create_run_idempotent(&request, key).await?,
163                None => client.create_run(&request).await?,
164            };
165            output::print_output(json_mode, &response, || {
166                output::runs_table(slice::from_ref(&response.data))
167            })?;
168        }
169        RunCommands::List {
170            status,
171            workflow,
172            created_by,
173            page,
174            per_page,
175        } => {
176            let filter = ListRunsFilter {
177                status: status.as_deref(),
178                workflow: workflow.as_deref(),
179                created_by: *created_by,
180                page: *page,
181                per_page: *per_page,
182                ..Default::default()
183            };
184            let response = client.list_runs_filtered(&filter).await?;
185            output::print_output(json_mode, &response, || output::runs_table(&response.data))?;
186        }
187        RunCommands::Get { id } => {
188            let response = client.get_run(*id).await?;
189            output::print_output(json_mode, &response, || {
190                output::run_detail_table(&response.data)
191            })?;
192
193            if !json_mode && !response.data.steps.is_empty() {
194                let mut out = std::io::stdout().lock();
195                writeln!(out)?;
196                writeln!(out, "Steps:")?;
197                writeln!(out, "{}", output::steps_table(&response.data.steps))?;
198            }
199        }
200        RunCommands::Cancel { id } => {
201            let response = client.cancel_run(*id).await?;
202            output::print_output(json_mode, &response, || {
203                output::runs_table(slice::from_ref(&response.data))
204            })?;
205        }
206        RunCommands::Approve { id } => {
207            let response = client.approve_run(*id).await?;
208            output::print_output(json_mode, &response, || {
209                output::runs_table(slice::from_ref(&response.data))
210            })?;
211        }
212        RunCommands::Retry { id } => {
213            let response = client.retry_run(*id).await?;
214            output::print_output(json_mode, &response, || {
215                output::runs_table(slice::from_ref(&response.data))
216            })?;
217        }
218    }
219    Ok(())
220}
221
222#[cfg(test)]
223mod tests {
224    use std::io::Write;
225
226    use tempfile::NamedTempFile;
227
228    use super::*;
229
230    #[test]
231    fn resolve_payload_none_returns_empty_object() {
232        let value = resolve_payload(None, None).unwrap();
233        assert!(value.is_object());
234        assert!(value.as_object().unwrap().is_empty());
235    }
236
237    #[test]
238    fn resolve_payload_inline_valid_json() {
239        let value = resolve_payload(Some(r#"{"key": "value"}"#), None).unwrap();
240        assert_eq!(value["key"], "value");
241    }
242
243    #[test]
244    fn resolve_payload_inline_invalid_json() {
245        let result = resolve_payload(Some("not json"), None);
246        assert!(result.is_err());
247        assert!(result.unwrap_err().to_string().contains("invalid JSON"));
248    }
249
250    #[test]
251    fn resolve_payload_file_valid() {
252        let mut tmp = NamedTempFile::new().unwrap();
253        write!(tmp, r#"{{"workflow": "test"}}"#).unwrap();
254        let path = tmp.path().to_path_buf();
255
256        let value = resolve_payload(None, Some(&path)).unwrap();
257        assert_eq!(value["workflow"], "test");
258    }
259
260    #[test]
261    fn resolve_payload_file_not_found() {
262        let path = PathBuf::from("/nonexistent/payload.json");
263        let result = resolve_payload(None, Some(&path));
264        assert!(result.is_err());
265        assert!(result.unwrap_err().to_string().contains("cannot read"));
266    }
267
268    #[test]
269    fn resolve_payload_file_invalid_json() {
270        let mut tmp = NamedTempFile::new().unwrap();
271        write!(tmp, "not valid json").unwrap();
272        let path = tmp.path().to_path_buf();
273
274        let result = resolve_payload(None, Some(&path));
275        assert!(result.is_err());
276        assert!(result.unwrap_err().to_string().contains("invalid JSON"));
277    }
278}