1use 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#[derive(Debug, Args)]
25pub struct RunArgs {
26 #[command(subcommand)]
28 pub command: RunCommands,
29}
30
31#[derive(Debug, Subcommand)]
33pub enum RunCommands {
34 Create {
36 workflow: String,
38 #[arg(long, group = "payload_source")]
40 payload: Option<String>,
41 #[arg(long, group = "payload_source")]
43 payload_file: Option<PathBuf>,
44 #[arg(long)]
47 max_retries: Option<u32>,
48 #[arg(long)]
54 idempotency_key: Option<String>,
55 #[arg(long = "max-cost", value_name = "USD")]
58 max_cost: Option<f64>,
59 },
60 List {
62 #[arg(long)]
64 status: Option<String>,
65 #[arg(long)]
67 workflow: Option<String>,
68 #[arg(long)]
72 created_by: Option<Uuid>,
73 #[arg(long)]
75 page: Option<u32>,
76 #[arg(long)]
78 per_page: Option<u32>,
79 },
80 Get {
82 id: Uuid,
84 },
85 Cancel {
87 id: Uuid,
89 },
90 Approve {
92 id: Uuid,
94 },
95 Reject {
97 id: Uuid,
99 },
100 Input {
102 id: Uuid,
104 step_id: Uuid,
106 #[arg(long, group = "value_source")]
108 value: Option<String>,
109 #[arg(long, group = "value_source")]
111 value_file: Option<PathBuf>,
112 },
113 RejectInput {
115 id: Uuid,
117 step_id: Uuid,
119 #[arg(long)]
121 reason: Option<String>,
122 },
123 Retry {
125 id: Uuid,
127 #[arg(long)]
130 force: bool,
131 },
132 Replay {
134 id: Uuid,
136 },
137 Watch {
139 id: Uuid,
141 #[arg(long)]
143 no_logs: bool,
144 #[arg(long, value_parser = parse_humantime)]
146 timeout: Option<Duration>,
147 },
148 Plan {
150 workflow: String,
152 #[arg(long, group = "plan_input_source")]
154 input: Option<String>,
155 #[arg(long, group = "plan_input_source")]
157 input_file: Option<PathBuf>,
158 #[arg(long)]
160 max_depth: Option<u32>,
161 #[arg(long)]
163 no_estimates: bool,
164 },
165 Diff {
167 run_a: Uuid,
169 run_b: Uuid,
171 },
172}
173
174fn parse_humantime(s: &str) -> Result<Duration, String> {
176 humantime::parse_duration(s).map_err(|e| e.to_string())
177}
178
179const TERMINAL_EVENTS: &[&str] = &["run_completed", "run_failed", "run_cancelled"];
181
182fn 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
195fn 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
215pub 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 .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 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 .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
385async 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
460async 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}